erebos 0.1.5 → 0.1.6
raw patch · 14 files changed
+361/−270 lines, 14 filesdep ~basedep ~template-haskellPVP ok
version bump matches the API change (PVP)
Dependency ranges changed: base, template-haskell
API changes (from Hackage documentation)
+ Erebos.Chatroom: chatroomMembers :: ChatroomState -> [ComposedIdentity]
+ Erebos.Chatroom: joinChatroom :: (MonadStorage m, MonadHead LocalState m, MonadError String m) => ChatroomState -> m ()
+ Erebos.Chatroom: joinChatroomByStateData :: (MonadStorage m, MonadHead LocalState m, MonadError String m) => Stored ChatroomStateData -> m ()
+ Erebos.Chatroom: leaveChatroom :: (MonadStorage m, MonadHead LocalState m, MonadError String m) => ChatroomState -> m ()
+ Erebos.Chatroom: leaveChatroomByStateData :: (MonadStorage m, MonadHead LocalState m, MonadError String m) => Stored ChatroomStateData -> m ()
+ Erebos.Network: sendManyToPeer :: (Service s, MonadIO m) => Peer -> [s] -> m ()
+ Erebos.Network: sendManyToPeerStored :: (Service s, MonadIO m) => Peer -> [Stored s] -> m ()
Files
- CHANGELOG.md +7/−0
- README.md +11/−0
- erebos.cabal +3/−4
- main/Main.hs +77/−32
- main/Test.hs +30/−6
- src/Erebos/Chatroom.hs +50/−15
- src/Erebos/Identity.hs +13/−20
- src/Erebos/Network.hs +77/−27
- src/Erebos/Network/Protocol.hs +1/−0
- src/Erebos/Network/ifaddrs.c +80/−4
- src/Erebos/Storage/Internal.hs +3/−3
- src/Erebos/Storage/Key.hs +4/−3
- src/Erebos/Storage/List.hs +0/−154
- src/Erebos/Storage/Merge.hs +5/−2
CHANGELOG.md view
@@ -1,5 +1,12 @@ # Revision history for erebos +## 0.1.6 -- 2024-08-12++* Fix sending multiple data responses in a stream+* Added `--storage`/`--memory-storage` command-line options+* Compatibility with GHC up to 9.10+* Local discovery with IPv6+ ## 0.1.5 -- 2024-07-16 * Public chatrooms for multiple participants
README.md view
@@ -121,6 +121,17 @@ : Create public unmoderated chatroom. Room name can be passed as command argument or entered interactively. +`/members` +: List members of the chatroom – usesers who sent any message or joined via the+`join` command.++`/join` +: Join chatroom without sending text message.++`/leave` +: Leave the chatroom. User will no longer be listed as a member and erebos tool+ will no longer collect message of this chatroom.+ ### Add contacts To ensure the identity of the contact and prevent man-in-the-middle attack,
erebos.cabal view
@@ -1,7 +1,7 @@ Cabal-Version: 3.0 Name: erebos-Version: 0.1.5+Version: 0.1.6 Synopsis: Decentralized messaging and synchronization Description: Library and simple CLI interface implementing the Erebos identity@@ -54,7 +54,7 @@ -Wno-error=unused-imports build-depends:- base >=4.13 && <4.20,+ base ^>= { 4.15, 4.16, 4.17, 4.18, 4.19, 4.20 }, default-extensions: DefaultSignatures@@ -116,7 +116,6 @@ Erebos.Storage.Internal other-modules: Erebos.Flow- Erebos.Storage.List Erebos.Storage.Platform Erebos.Util @@ -195,7 +194,7 @@ mtl, network, process >=1.6 && <1.7,- template-haskell >=2.17 && <2.22,+ template-haskell ^>= { 2.17, 2.18, 2.19, 2.20, 2.21, 2.22 }, text, time, transformers >= 0.5 && <0.7,
main/Main.hs view
@@ -61,12 +61,17 @@ data Options = Options { optServer :: ServerOptions , optServices :: [ServiceOption]+ , optStorage :: StorageOption , optChatroomAutoSubscribe :: Maybe Int , optDmBotEcho :: Maybe Text , optShowHelp :: Bool , optShowVersion :: Bool } +data StorageOption = DefaultStorage+ | FilesystemStorage FilePath+ | MemoryStorage+ data ServiceOption = ServiceOption { soptName :: String , soptService :: SomeService@@ -78,6 +83,7 @@ defaultOptions = Options { optServer = defaultServerOptions , optServices = availableServices+ , optStorage = DefaultStorage , optChatroomAutoSubscribe = Nothing , optDmBotEcho = Nothing , optShowHelp = False@@ -110,6 +116,12 @@ , Option ['s'] ["silent"] (NoArg (so $ \opts -> opts { serverLocalDiscovery = False })) "do not send announce packets for local discovery"+ , Option [] [ "storage" ]+ (ReqArg (\path -> \opts -> opts { optStorage = FilesystemStorage path }) "<path>")+ "use storage in <path>"+ , Option [] [ "memory-storage" ]+ (NoArg (\opts -> opts { optStorage = MemoryStorage }))+ "use memory storage" , Option [] ["chatroom-auto-subscribe"] (ReqArg (\count -> \opts -> opts { optChatroomAutoSubscribe = Just (read count) }) "<count>") "automatically subscribe for up to <count> chatrooms"@@ -142,8 +154,20 @@ main :: IO () main = do- st <- liftIO $ openStorage . fromMaybe "./.erebos" =<< lookupEnv "EREBOS_DIR"- getArgs >>= \case+ (opts, args) <- (getOpt RequireOrder (options ++ servicesOptions) <$> getArgs) >>= \case+ (o, args, []) -> do+ return (foldl (flip id) defaultOptions o, args)+ (_, _, errs) -> do+ progName <- getProgName+ hPutStrLn stderr $ concat errs <> "Try `" <> progName <> " --help' for more information."+ exitFailure++ st <- liftIO $ case optStorage opts of+ DefaultStorage -> openStorage . fromMaybe "./.erebos" =<< lookupEnv "EREBOS_DIR"+ FilesystemStorage path -> openStorage path+ MemoryStorage -> memoryStorage++ case args of ["cat-file", sref] -> do readRef st (BC.pack sref) >>= \case Nothing -> error "ref does not exist"@@ -159,7 +183,7 @@ forM_ (signedSignature signed) $ \sig -> do putStr $ "SIG " BC.putStrLn $ showRef $ storedRef $ sigKey $ fromStored sig- "identity" -> case validateIdentityF (wrappedLoad <$> refs) of+ "identity" -> case validateExtendedIdentityF (wrappedLoad <$> refs) of Just identity -> do let disp :: Identity m -> IO () disp idt = do@@ -169,7 +193,7 @@ case idOwner idt of Nothing -> return () Just owner -> do- mapM_ (putStrLn . ("OWNER " ++) . BC.unpack . showRefDigest . refDigest . storedRef) $ idDataF owner+ mapM_ (putStrLn . ("OWNER " ++) . BC.unpack . showRefDigest . refDigest . storedRef) $ idExtDataF owner disp owner disp identity Nothing -> putStrLn $ "Identity verification failed"@@ -193,34 +217,32 @@ ["test"] -> runTestTool st - args -> case getOpt Permute (options ++ servicesOptions) args of- (o, [], []) -> do- let opts = foldl (flip id) defaultOptions o- header = "Usage: erebos [OPTION...]"- serviceDesc ServiceOption {..} = padService (" " <> soptName) <> soptDescription+ [] -> do+ let header = "Usage: erebos [OPTION...]"+ serviceDesc ServiceOption {..} = padService (" " <> soptName) <> soptDescription - padTo n str = str <> replicate (n - length str) ' '- padOpt = padTo 37- padService = padTo 16+ padTo n str = str <> replicate (n - length str) ' '+ padOpt = padTo 37+ padService = padTo 16 - if | optShowHelp opts -> putStr $ usageInfo header options <> unlines- (- [ padOpt " --enable-<service>" <> "enable network service <service>"- , padOpt " --disable-<service>" <> "disable network service <service>"- , padOpt " --enable-all" <> "enable all network services"- , padOpt " --disable-all" <> "disable all network services"- , ""- , "Available network services:"- ] ++ map serviceDesc availableServices- )- | optShowVersion opts -> putStrLn versionLine- | otherwise -> interactiveLoop st opts- (_, _, errs) -> do- progName <- getProgName- hPutStrLn stderr $ concat errs <> "Try `" <> progName <> " --help' for more information."- exitFailure+ if | optShowHelp opts -> putStr $ usageInfo header options <> unlines+ (+ [ padOpt " --enable-<service>" <> "enable network service <service>"+ , padOpt " --disable-<service>" <> "disable network service <service>"+ , padOpt " --enable-all" <> "enable all network services"+ , padOpt " --disable-all" <> "disable all network services"+ , ""+ , "Available network services:"+ ] ++ map serviceDesc availableServices+ )+ | optShowVersion opts -> putStrLn versionLine+ | otherwise -> interactiveLoop st opts + (cmdname : _) -> do+ hPutStrLn stderr $ "Unknown command `" <> cmdname <> "'"+ exitFailure + inputSettings :: Settings IO inputSettings = setComplete commandCompletion $ defaultSettings @@ -231,8 +253,10 @@ tui <- haveTerminalUI extPrint <- getExternalPrint- let extPrintLn str = extPrint $ case reverse str of ('\n':_) -> str- _ -> str ++ "\n";+ let extPrintLn str = do+ let str' = case reverse str of ('\n':_) -> str+ _ -> str ++ "\n";+ extPrint $! str' -- evaluate str before calling extPrint to avoid blinking let getInputLinesTui eprompt = do prompt <- case eprompt of@@ -428,6 +452,11 @@ SelectedPeer peer -> return peer _ -> throwError "no peer selected" +getSelectedChatroom :: CommandM ChatroomState+getSelectedChatroom = gets csContext >>= \case+ SelectedChatroom rstate -> return rstate+ _ -> throwError "no chatroom selected"+ getSelectedConversation :: CommandM Conversation getSelectedConversation = gets csContext >>= \case SelectedPeer peer -> peerIdentity peer >>= \case@@ -472,6 +501,9 @@ , ("ice-connect", cmdIceConnect) , ("ice-send", cmdIceSend) #endif+ , ("join", cmdJoin)+ , ("leave", cmdLeave)+ , ("members", cmdMembers) , ("select", cmdSelectContext) , ("quit", cmdQuit) ]@@ -524,6 +556,19 @@ PeerIdentityFull pid -> T.unpack $ displayIdentity pid in name ++ " [" ++ show paddr ++ "]" +cmdJoin :: Command+cmdJoin = joinChatroom =<< getSelectedChatroom++cmdLeave :: Command+cmdLeave = leaveChatroom =<< getSelectedChatroom++cmdMembers :: Command+cmdMembers = do+ Just room <- findChatroomByStateData . head . roomStateData =<< getSelectedChatroom+ forM_ (chatroomMembers room) $ \x -> do+ liftIO $ putStrLn $ maybe "<unnamed>" T.unpack $ idName x++ cmdSelectContext :: Command cmdSelectContext = do n <- read <$> asks ciLine@@ -629,8 +674,8 @@ [ maybe "<unnamed>" T.unpack $ roomName =<< cmsgRoom msg , formatTime defaultTimeLocale " [%H:%M] " $ utcToLocalTime tzone $ zonedTimeToUTC $ cmsgTime msg , maybe "<unnamed>" T.unpack $ idName $ cmsgFrom msg- , ": "- , maybe "<no message>" T.unpack $ cmsgText msg+ , if cmsgLeave msg then " left" else ""+ , maybe (if cmsgLeave msg then "" else " joined") ((": " ++) . T.unpack) $ cmsgText msg ] modifyMVar_ subscribedNumVar $ return . (if roomStateSubscribe rstate then (+ 1) else id)
main/Test.hs view
@@ -97,7 +97,7 @@ Nothing -> return () runExceptT (evalStateT testLoop initTestState) >>= \case- Left x -> hPutStrLn stderr x+ Left x -> B.hPutStr stderr $ (`BC.snoc` '\n') $ BC.pack x Right () -> return () getLineMb :: MonadIO m => m (Maybe Text)@@ -121,7 +121,7 @@ outLine mvar line = do evaluate $ foldl' (flip seq) () line withMVar mvar $ \() -> do- putStrLn line+ B.putStr $ (`BC.snoc` '\n') $ BC.pack line hFlush stdout cmdOut :: String -> Command@@ -283,6 +283,9 @@ , ("chatroom-set-name", cmdChatroomSetName) , ("chatroom-subscribe", cmdChatroomSubscribe) , ("chatroom-unsubscribe", cmdChatroomUnsubscribe)+ , ("chatroom-members", cmdChatroomMembers)+ , ("chatroom-join", cmdChatroomJoin)+ , ("chatroom-leave", cmdChatroomLeave) , ("chatroom-message-send", cmdChatroomMessageSend) ] @@ -428,7 +431,7 @@ h <- getOrLoadHead rsPeers <- liftIO $ newMVar (1, [])- rsServer <- liftIO $ startServer defaultServerOptions h (hPutStrLn stderr)+ rsServer <- liftIO $ startServer defaultServerOptions h (B.hPutStr stderr . (`BC.snoc` '\n') . BC.pack) [ someServiceAttr $ pairingAttributes (Proxy @AttachService) out rsPeers "attach" , someServiceAttr $ pairingAttributes (Proxy @ContactService) out rsPeers "contact" , someServiceAttr $ directMessageAttributes out@@ -501,11 +504,11 @@ cmdTestMessageSend :: Command cmdTestMessageSend = do- [spidx, tref] <- asks tiParams+ spidx : trefs <- asks tiParams st <- asks tiStorage- Just ref <- liftIO $ readRef st (encodeUtf8 tref)+ Just refs <- liftIO $ fmap sequence $ mapM (readRef st . encodeUtf8) trefs peer <- getPeer spidx- sendToPeer peer $ TestMessage $ wrappedLoad ref+ sendManyToPeer peer $ map (TestMessage . wrappedLoad) refs cmdOut "test-message-send done" cmdSharedStateGet :: Command@@ -732,6 +735,7 @@ , [ show . refDigest . storedRef . head . filterAncestors . concatMap storedRoots . toComponents $ room ] , [ "room", maybe "<unnamed>" T.unpack $ roomName =<< cmsgRoom msg ] , [ "from", maybe "<unnamed>" T.unpack $ idName $ cmsgFrom msg ]+ , if cmsgLeave msg then [ "leave" ] else [] , maybe [] (("text":) . (:[]) . T.unpack) $ cmsgText msg ] @@ -753,6 +757,26 @@ [ cid ] <- asks tiParams to <- getChatroomStateData cid void $ chatroomSetSubscribe to False++cmdChatroomMembers :: Command+cmdChatroomMembers = do+ [ cid ] <- asks tiParams+ Just chatroom <- findChatroomByStateData =<< getChatroomStateData cid+ forM_ (chatroomMembers chatroom) $ \user -> do+ cmdOut $ unwords [ "chatroom-members-item", maybe "<unnamed>" T.unpack $ idName user ]+ cmdOut "chatroom-members-done"++cmdChatroomJoin :: Command+cmdChatroomJoin = do+ [ cid ] <- asks tiParams+ joinChatroomByStateData =<< getChatroomStateData cid+ cmdOut "chatroom-join-done"++cmdChatroomLeave :: Command+cmdChatroomLeave = do+ [ cid ] <- asks tiParams+ leaveChatroomByStateData =<< getChatroomStateData cid+ cmdOut "chatroom-leave-done" cmdChatroomMessageSend :: Command cmdChatroomMessageSend = do
src/Erebos/Chatroom.hs view
@@ -11,6 +11,9 @@ findChatroomByRoomData, findChatroomByStateData, chatroomSetSubscribe,+ chatroomMembers,+ joinChatroom, joinChatroomByStateData,+ leaveChatroom, leaveChatroomByStateData, getMessagesSinceState, ChatroomSetChange(..),@@ -33,6 +36,8 @@ import Data.Bool import Data.Either+import Data.Foldable+import Data.Function import Data.IORef import Data.List import Data.Maybe@@ -180,23 +185,23 @@ sendChatroomMessageByStateData :: (MonadStorage m, MonadHead LocalState m, MonadError String m) => Stored ChatroomStateData -> Text -> m ()-sendChatroomMessageByStateData lookupData msg = void $ findAndUpdateChatroomState $ \cstate -> do+sendChatroomMessageByStateData lookupData msg = sendRawChatroomMessageByStateData lookupData Nothing (Just msg) False++sendRawChatroomMessageByStateData+ :: (MonadStorage m, MonadHead LocalState m, MonadError String m)+ => Stored ChatroomStateData -> Maybe (Stored (Signed ChatMessageData)) -> Maybe Text -> Bool -> m ()+sendRawChatroomMessageByStateData lookupData mdReplyTo mdText mdLeave = void $ findAndUpdateChatroomState $ \cstate -> do guard $ any (lookupData `precedesOrEquals`) $ roomStateData cstate Just $ do- self <- finalOwner . localIdentity . fromStored <$> getLocalHead- secret <- loadKey $ idKeyMessage self- time <- liftIO getZonedTime- mdata <- mstore =<< sign secret =<< mstore ChatMessageData- { mdPrev = roomStateMessageData cstate- , mdRoom = if null (roomStateMessageData cstate)- then maybe [] roomData (roomStateRoom cstate)- else []- , mdFrom = self- , mdReplyTo = Nothing- , mdTime = time- , mdText = Just msg- , mdLeave = False- }+ mdFrom <- finalOwner . localIdentity . fromStored <$> getLocalHead+ secret <- loadKey $ idKeyMessage mdFrom+ mdTime <- liftIO getZonedTime+ let mdPrev = roomStateMessageData cstate+ mdRoom = if null (roomStateMessageData cstate)+ then maybe [] roomData (roomStateRoom cstate)+ else []++ mdata <- mstore =<< sign secret =<< mstore ChatMessageData {..} mergeSorted . (:[]) <$> mstore ChatroomStateData { rsdPrev = roomStateData cstate , rsdRoom = []@@ -340,6 +345,36 @@ , rsdSubscribe = Just subscribe , rsdMessages = [] }++chatroomMembers :: ChatroomState -> [ ComposedIdentity ]+chatroomMembers ChatroomState {..} =+ map (mdFrom . fromSigned . head) $+ filter (any $ not . mdLeave . fromSigned) $ -- keep only users that hasn't left+ map (filterAncestors . map snd) $ -- gather message data per each identity and filter ancestors+ groupBy ((==) `on` fst) $ -- group on identity root+ sortBy (comparing fst) $ -- sort by first root of identity data+ map (\x -> ( head . filterAncestors . concatMap storedRoots . idDataF . mdFrom . fromSigned $ x, x )) $+ toList $ ancestors $ roomStateMessageData++joinChatroom+ :: (MonadStorage m, MonadHead LocalState m, MonadError String m)+ => ChatroomState -> m ()+joinChatroom rstate = joinChatroomByStateData (head $ roomStateData rstate)++joinChatroomByStateData+ :: (MonadStorage m, MonadHead LocalState m, MonadError String m)+ => Stored ChatroomStateData -> m ()+joinChatroomByStateData lookupData = sendRawChatroomMessageByStateData lookupData Nothing Nothing False++leaveChatroom+ :: (MonadStorage m, MonadHead LocalState m, MonadError String m)+ => ChatroomState -> m ()+leaveChatroom rstate = leaveChatroomByStateData (head $ roomStateData rstate)++leaveChatroomByStateData+ :: (MonadStorage m, MonadHead LocalState m, MonadError String m)+ => Stored ChatroomStateData -> m ()+leaveChatroomByStateData lookupData = sendRawChatroomMessageByStateData lookupData Nothing Nothing True getMessagesSinceState :: ChatroomState -> ChatroomState -> [ChatMessage] getMessagesSinceState cur old = threadToListSince (roomStateMessageData old) (roomStateMessageData cur)
src/Erebos/Identity.hs view
@@ -35,7 +35,6 @@ import Data.Function import Data.List import Data.Maybe-import Data.Ord import Data.Set (Set) import qualified Data.Set as S import Data.Text (Text)@@ -304,25 +303,18 @@ throwError "signature verification failed" lookupProperty :: forall a m. Foldable m => (ExtendedIdentityData -> Maybe a) -> m (Stored (Signed ExtendedIdentityData)) -> Maybe a-lookupProperty sel topHeads = findResult filteredLayers- where findPropHeads :: Stored (Signed ExtendedIdentityData) -> [(Stored (Signed ExtendedIdentityData), a)]- findPropHeads sobj | Just x <- sel $ fromSigned sobj = [(sobj, x)]- | otherwise = findPropHeads =<< (eiddPrev $ fromSigned sobj)-- propHeads :: [(Stored (Signed ExtendedIdentityData), a)]- propHeads = findPropHeads =<< toList topHeads-- historyLayers :: [Set (Stored (Signed ExtendedIdentityData))]- historyLayers = generations $ map fst propHeads+lookupProperty sel topHeads = findResult propHeads+ where+ findPropHeads :: Stored (Signed ExtendedIdentityData) -> [ Stored (Signed ExtendedIdentityData) ]+ findPropHeads sobj | Just _ <- sel $ fromSigned sobj = [ sobj ]+ | otherwise = findPropHeads =<< (eiddPrev $ fromSigned sobj) - filteredLayers :: [[(Stored (Signed ExtendedIdentityData), a)]]- filteredLayers = scanl (\cur obsolete -> filter ((`S.notMember` obsolete) . fst) cur) propHeads historyLayers+ propHeads :: [ Stored (Signed ExtendedIdentityData) ]+ propHeads = filterAncestors $ findPropHeads =<< toList topHeads - findResult ([(_, x)] : _) = Just x- findResult ([] : _) = Nothing- findResult [] = Nothing- findResult [xs] = Just $ snd $ minimumBy (comparing fst) xs- findResult (_:rest) = findResult rest+ findResult :: [ Stored (Signed ExtendedIdentityData) ] -> Maybe a+ findResult [] = Nothing+ findResult xs = sel $ fromSigned $ minimum xs mergeIdentity :: (MonadStorage m, MonadError String m, MonadIO m) => Identity f -> m UnifiedIdentity mergeIdentity idt | Just idt' <- toUnifiedIdentity idt = return idt'@@ -385,8 +377,9 @@ updateOwners _ orig@Identity { idOwner_ = Nothing } = orig sameIdentity :: (Foldable m, Foldable m') => Identity m -> Identity m' -> Bool-sameIdentity x y = not $ S.null $ S.intersection (refset x) (refset y)- where refset idt = foldr S.insert (ancestors $ toList $ idDataF idt) (idDataF idt)+sameIdentity x y = intersectsSorted (roots x) (roots y)+ where+ roots idt = uniq $ sort $ concatMap storedRoots $ toList $ idDataF idt unfoldOwners :: (Foldable m) => Identity m -> [ComposedIdentity]
src/Erebos/Network.hs view
@@ -19,7 +19,9 @@ #endif dropPeer, isPeerDropped,- sendToPeer, sendToPeerStored, sendToPeerWith,+ sendToPeer, sendManyToPeer,+ sendToPeerStored, sendManyToPeerStored,+ sendToPeerWith, runPeerService, discoveryPort,@@ -52,6 +54,9 @@ import Network.Socket hiding (ControlMessage) import qualified Network.Socket.ByteString as S +import Foreign.C.Types+import Foreign.Marshal.Alloc+ import Erebos.Channel #ifdef ENABLE_ICE_SUPPORT import Erebos.ICE@@ -69,6 +74,9 @@ discoveryPort :: PortNumber discoveryPort = 29665 +discoveryMulticastGroup :: HostAddress6+discoveryMulticastGroup = tupleToHostAddress6 (0xff12, 0xb6a4, 0x6b1f, 0x0969, 0xcaee, 0xacc2, 0x5c93, 0x73e1) -- ff12:b6a4:6b1f:969:caee:acc2:5c93:73e1+ announceIntervalSeconds :: Int announceIntervalSeconds = 60 @@ -247,8 +255,6 @@ either (atomically . logd) return =<< runExceptT =<< atomically (readTQueue serverIOActions) - broadcastAddreses <- getBroadcastAddresses discoveryPort- let open addr = do sock <- socket (addrFamily addr) (addrSocketType addr) (addrProtocol addr) putMVar serverSocket sock@@ -259,9 +265,14 @@ return sock loop sock = do- when (serverLocalDiscovery opt) $ forkServerThread server $ forever $ do- atomically $ writeFlowBulk serverControlFlow $ map (SendAnnounce . DatagramAddress) broadcastAddreses- threadDelay $ announceIntervalSeconds * 1000 * 1000+ when (serverLocalDiscovery opt) $ forkServerThread server $ do+ announceAddreses <- fmap concat $ sequence $+ [ map (SockAddrInet6 discoveryPort 0 discoveryMulticastGroup) <$> joinMulticast sock+ , getBroadcastAddresses discoveryPort+ ]+ forever $ do+ atomically $ writeFlowBulk serverControlFlow $ map (SendAnnounce . DatagramAddress) announceAddreses+ threadDelay $ announceIntervalSeconds * 1000 * 1000 let announceUpdate identity = do st <- derivePartialStorage serverStorage@@ -422,12 +433,18 @@ runPacketHandler :: Bool -> Peer -> PacketHandler () -> STM () runPacketHandler secure peer@Peer {..} act = do let logd = writeTQueue $ serverErrorLog peerServer_- runExceptT (flip execStateT (PacketHandlerState peer [] [] [] False) $ unPacketHandler act) >>= \case+ runExceptT (flip execStateT (PacketHandlerState peer [] [] [] Nothing False) $ unPacketHandler act) >>= \case Left err -> do logd $ "Error in handling packet from " ++ show peerAddress ++ ": " ++ err Right ph -> do when (not $ null $ phHead ph) $ do- let packet = TransportPacket (TransportHeader $ phHead ph) (phBody ph)+ body <- case phBodyStream ph of+ Nothing -> return $ phBody ph+ Just stream -> do+ writeTQueue (serverIOActions peerServer_) $ void $ liftIO $ forkIO $ do+ writeByteStringToStream stream $ BL.concat $ map lazyLoadBytes $ phBody ph+ return []+ let packet = TransportPacket (TransportHeader $ phHead ph) body secreq = case (secure, phPlaintextReply ph) of (True, _) -> EncryptedOnly (False, False) -> PlaintextAllowed@@ -451,6 +468,7 @@ , phHead :: [TransportHeaderItem] , phAckedBy :: [TransportHeaderItem] , phBody :: [Ref]+ , phBodyStream :: Maybe RawStreamWriter , phPlaintextReply :: Bool } @@ -463,6 +481,14 @@ addBody :: Ref -> PacketHandler () addBody r = modify $ \ph -> ph { phBody = r `appendDistinct` phBody ph } +sendBodyAsStream :: PacketHandler ()+sendBodyAsStream = do+ gets phBodyStream >>= \case+ Nothing -> do+ stream <- openStream+ modify $ \ph -> ph { phBodyStream = Just stream }+ Just _ -> return ()+ keepPlaintextReply :: PacketHandler () keepPlaintextReply = modify $ \ph -> ph { phPlaintextReply = True } @@ -518,9 +544,13 @@ liftSTM $ finalizedChannel peer ch identity _ -> return () - Rejected dgst -> do- logd $ "rejected by peer: " ++ show dgst+ Rejected dgst+ | peerRequest : _ <- mapMaybe (\case TrChannelRequest d -> Just d; _ -> Nothing) headers+ , peerRequest < dgst+ -> return () -- Our request was rejected due to lower priority + | otherwise -> logd $ "rejected by peer: " ++ show dgst+ DataRequest dgst | secure || dgst `elem` plaintextRefs -> do Right mref <- liftSTM $ unsafeIOToSTM $@@ -533,15 +563,11 @@ -- otherwise lost the channel, so keep the reply plaintext as well. when (not secure) keepPlaintextReply - let bytes = lazyLoadBytes mref -- TODO: MTU- if (secure && BL.length bytes > 500)- then do- stream <- openStream- liftSTM $ writeTQueue (serverIOActions server) $ void $ liftIO $ forkIO $ do- writeByteStringToStream stream bytes- else do- addBody $ mref+ when (secure && BL.length (lazyLoadBytes mref) > 500)+ sendBodyAsStream++ addBody $ mref | otherwise -> do logd $ "unauthorized data request for " ++ show dgst addHeader $ Rejected dgst@@ -594,9 +620,15 @@ ChannelCookieWait {} -> return () ChannelCookieReceived {} -> process ChannelCookieConfirmed {} -> process- ChannelOurRequest our | dgst < refDigest (storedRef our) -> process- | otherwise -> reject- ChannelPeerRequest {} -> process+ ChannelOurRequest our+ | dgst < refDigest (storedRef our) -> process+ | otherwise -> do+ -- Reject peer channel request with lower priority+ addHeader $ TrChannelRequest $ refDigest $ storedRef our+ reject+ ChannelPeerRequest prev+ | dgst == wrDigest prev -> addHeader $ Acknowledged dgst+ | otherwise -> process ChannelOurAccept {} -> reject ChannelEstablished {} -> process ChannelClosed {} -> return ()@@ -648,12 +680,14 @@ [ TrChannelRequest reqref , AnnounceSelf $ refDigest $ storedRef $ idData identity ]+ let sendChannelRequest = do+ sendToPeerPlain peer [ Acknowledged reqref, Rejected reqref ] $+ TransportPacket (TransportHeader hitems) [storedRef req]+ setPeerChannel peer $ ChannelOurRequest req liftIO $ atomically $ do getPeerChannel peer >>= \case- ChannelCookieConfirmed -> do- sendToPeerPlain peer [ Acknowledged reqref, Rejected reqref ] $- TransportPacket (TransportHeader hitems) [storedRef req]- setPeerChannel peer $ ChannelOurRequest req+ ChannelCookieReceived -> sendChannelRequest+ ChannelCookieConfirmed -> sendChannelRequest _ -> return () handleChannelRequest :: Peer -> UnifiedIdentity -> Ref -> WaitingRefCallback@@ -807,11 +841,17 @@ _ -> return False sendToPeer :: (Service s, MonadIO m) => Peer -> s -> m ()-sendToPeer peer packet = sendToPeerList peer [ServiceReply (Left packet) True]+sendToPeer peer = sendManyToPeer peer . (: []) +sendManyToPeer :: (Service s, MonadIO m) => Peer -> [ s ] -> m ()+sendManyToPeer peer = sendToPeerList peer . map (\part -> ServiceReply (Left part) True)+ sendToPeerStored :: (Service s, MonadIO m) => Peer -> Stored s -> m ()-sendToPeerStored peer spacket = sendToPeerList peer [ServiceReply (Right spacket) True]+sendToPeerStored peer = sendManyToPeerStored peer . (: []) +sendManyToPeerStored :: (Service s, MonadIO m) => Peer -> [ Stored s ] -> m ()+sendManyToPeerStored peer = sendToPeerList peer . map (\part -> ServiceReply (Right part) True)+ sendToPeerList :: (Service s, MonadIO m) => Peer -> [ServiceReply s] -> m () sendToPeerList peer parts = do let st = peerStorage peer@@ -913,8 +953,18 @@ logd $ "unhandled service '" ++ show (toUUID svc) ++ "'" +foreign import ccall unsafe "Network/ifaddrs.h join_multicast" cJoinMulticast :: CInt -> Ptr CSize -> IO (Ptr Word32) foreign import ccall unsafe "Network/ifaddrs.h broadcast_addresses" cBroadcastAddresses :: IO (Ptr Word32) foreign import ccall unsafe "stdlib.h free" cFree :: Ptr Word32 -> IO ()++joinMulticast :: Socket -> IO [ Word32 ]+joinMulticast sock =+ withFdSocket sock $ \fd ->+ alloca $ \pcount -> do+ ptr <- cJoinMulticast fd pcount+ count <- fromIntegral <$> peek pcount+ forM [ 0 .. count - 1 ] $ \i ->+ peekElemOff ptr i getBroadcastAddresses :: PortNumber -> IO [SockAddr] getBroadcastAddresses port = do
src/Erebos/Network/Protocol.hs view
@@ -891,6 +891,7 @@ now <- readTVar gNowVar if next <= now then do+ writeTVar cNextKeepAlive Nothing identity <- fst <$> readTVar gIdentity let header = TransportHeader [ AnnounceSelf $ refDigest $ storedRef $ idData identity ] writeTQueue cSecureOutQueue (EncryptedOnly, TransportPacket header [], [])
src/Erebos/Network/ifaddrs.c view
@@ -1,13 +1,89 @@ #include "ifaddrs.h" -#ifndef _WIN32+#include <errno.h>+#include <stdbool.h>+#include <stdio.h>+#include <stdlib.h>+#include <string.h> +#ifndef _WIN32 #include <arpa/inet.h>-#include <ifaddrs.h> #include <net/if.h>-#include <stdlib.h>-#include <sys/types.h>+#include <ifaddrs.h> #include <endian.h>+#include <sys/types.h>+#include <sys/socket.h>+#else+#include <winsock2.h>+#include <ws2ipdef.h>+#include <ws2tcpip.h>+#endif++#define DISCOVERY_MULTICAST_GROUP "ff12:b6a4:6b1f:969:caee:acc2:5c93:73e1"++uint32_t * join_multicast(int fd, size_t * count)+{+ size_t capacity = 16;+ *count = 0;+ uint32_t * interfaces = malloc(sizeof(uint32_t) * capacity);++#ifdef _WIN32+ interfaces[0] = 0;+ *count = 1;+#else+ struct ifaddrs * addrs;+ if (getifaddrs(&addrs) < 0)+ return 0;++ for (struct ifaddrs * ifa = addrs; ifa; ifa = ifa->ifa_next) {+ if (ifa->ifa_addr && ifa->ifa_addr->sa_family == AF_INET6 &&+ !(ifa->ifa_flags & IFF_LOOPBACK)) {+ int idx = if_nametoindex(ifa->ifa_name);++ bool seen = false;+ for (size_t i = 0; i < *count; i++) {+ if (interfaces[i] == idx) {+ seen = true;+ break;+ }+ }+ if (seen)+ continue;++ if (*count + 1 >= capacity) {+ capacity *= 2;+ uint32_t * nret = realloc(interfaces, sizeof(uint32_t) * capacity);+ if (nret) {+ interfaces = nret;+ } else {+ free(interfaces);+ *count = 0;+ return NULL;+ }+ }++ interfaces[*count] = idx;+ (*count)++;+ }+ }++ freeifaddrs(addrs);+#endif++ for (size_t i = 0; i < *count; i++) {+ struct ipv6_mreq group;+ group.ipv6mr_interface = interfaces[i];+ inet_pton(AF_INET6, DISCOVERY_MULTICAST_GROUP, &group.ipv6mr_multiaddr);+ int ret = setsockopt(fd, IPPROTO_IPV6, IPV6_ADD_MEMBERSHIP,+ (const void *) &group, sizeof(group));+ if (ret < 0)+ fprintf(stderr, "IPV6_ADD_MEMBERSHIP failed: %s\n", strerror(errno));+ }++ return interfaces;+}++#ifndef _WIN32 uint32_t * broadcast_addresses(void) {
src/Erebos/Storage/Internal.hs view
@@ -241,7 +241,7 @@ doesFileExist file >>= \case True -> removeFile locked False -> do BL.hPut h content- hFlush h+ hClose h renameFile locked file where locked = file ++ ".lock" @@ -254,13 +254,13 @@ removeFile locked return $ Left $ Just current (Nothing, False) -> do B.hPut h content- hFlush h+ hClose h renameFile locked file return $ Right () (Just expected, True) -> do current <- B.readFile file if current == expected then do B.hPut h content- hFlush h+ hClose h renameFile locked file return $ return () else do removeFile locked
src/Erebos/Storage/Key.hs view
@@ -80,6 +80,7 @@ return M.empty (StorageMemory { memKeys = fromKeys }, StorageMemory { memKeys = toKeys }) -> do- modifyMVar_ fromKeys $ \fkeys -> do- modifyMVar_ toKeys $ return . M.union fkeys- return M.empty+ when (fromKeys /= toKeys) $ do+ modifyMVar_ fromKeys $ \fkeys -> do+ modifyMVar_ toKeys $ return . M.union fkeys+ return M.empty
− src/Erebos/Storage/List.hs
@@ -1,154 +0,0 @@-module Erebos.Storage.List (- StoredList,- emptySList, fromSList, storedFromSList,- slistAdd, slistAddS,- -- TODO slistInsert, slistInsertS,- slistRemove, slistReplace, slistReplaceS,- -- TODO mapFromSList, updateOld,-- -- TODO StoreUpdate(..),- -- TODO withStoredListItem, withStoredListItemS,-) where--import Data.List-import Data.Maybe-import qualified Data.Set as S--import Erebos.Storage-import Erebos.Storage.Internal-import Erebos.Storage.Merge--data List a = ListNil- | ListItem { listPrev :: [StoredList a]- , listItem :: Maybe (Stored a)- , listRemove :: Maybe (Stored (List a))- }--type StoredList a = Stored (List a)--instance Storable a => Storable (List a) where- store' ListNil = storeZero- store' x@ListItem {} = storeRec $ do- mapM_ (storeRef "PREV") $ listPrev x- mapM_ (storeRef "item") $ listItem x- mapM_ (storeRef "remove") $ listRemove x-- load' = loadCurrentObject >>= \case- ZeroObject -> return ListNil- _ -> loadRec $ ListItem <$> loadRefs "PREV"- <*> loadMbRef "item"- <*> loadMbRef "remove"--instance Storable a => ZeroStorable (List a) where- fromZero _ = ListNil---emptySList :: Storable a => Storage -> IO (StoredList a)-emptySList st = wrappedStore st ListNil--groupsFromSLists :: forall a. Storable a => StoredList a -> [[Stored a]]-groupsFromSLists = helperSelect S.empty . (:[])- where- helperSelect :: S.Set (StoredList a) -> [StoredList a] -> [[Stored a]]- helperSelect rs xxs | x:xs <- sort $ filterRemoved rs xxs = helper rs x xs- | otherwise = []-- helper :: S.Set (StoredList a) -> StoredList a -> [StoredList a] -> [[Stored a]]- helper rs x xs- | ListNil <- fromStored x- = []-- | Just rm <- listRemove (fromStored x)- , ans <- ancestors [x]- , (other, collision) <- partition (S.null . S.intersection ans . ancestors . (:[])) xs- , cont <- helperSelect (rs `S.union` ancestors [rm]) $ concatMap (listPrev . fromStored) (x : collision) ++ other- = case catMaybes $ map (listItem . fromStored) (x : collision) of- [] -> cont- xis -> xis : cont-- | otherwise = case listItem (fromStored x) of- Nothing -> helperSelect rs $ listPrev (fromStored x) ++ xs- Just xi -> [xi] : (helperSelect rs $ listPrev (fromStored x) ++ xs)-- filterRemoved :: S.Set (StoredList a) -> [StoredList a] -> [StoredList a]- filterRemoved rs = filter (S.null . S.intersection rs . ancestors . (:[]))--fromSList :: Mergeable a => StoredList (Component a) -> [a]-fromSList = map merge . groupsFromSLists--storedFromSList :: (Mergeable a, Storable a) => StoredList (Component a) -> IO [Stored a]-storedFromSList = mapM storeMerge . groupsFromSLists--slistAdd :: Storable a => a -> StoredList a -> IO (StoredList a)-slistAdd x prev@(Stored (Ref st _) _) = do- sx <- wrappedStore st x- slistAddS sx prev--slistAddS :: Storable a => Stored a -> StoredList a -> IO (StoredList a)-slistAddS sx prev@(Stored (Ref st _) _) = wrappedStore st (ListItem [prev] (Just sx) Nothing)--{- TODO-slistInsert :: Storable a => Stored a -> a -> StoredList a -> IO (StoredList a)-slistInsert after x prev@(Stored (Ref st _) _) = do- sx <- wrappedStore st x- slistInsertS after sx prev--slistInsertS :: Storable a => Stored a -> Stored a -> StoredList a -> IO (StoredList a)-slistInsertS after sx prev@(Stored (Ref st _) _) = wrappedStore st $ ListItem Nothing (findSListRef after prev) (Just sx) prev--}--slistRemove :: Storable a => Stored a -> StoredList a -> IO (StoredList a)-slistRemove rm prev@(Stored (Ref st _) _) = wrappedStore st $ ListItem [prev] Nothing (findSListRef rm prev)--slistReplace :: Storable a => Stored a -> a -> StoredList a -> IO (StoredList a)-slistReplace rm x prev@(Stored (Ref st _) _) = do- sx <- wrappedStore st x- slistReplaceS rm sx prev--slistReplaceS :: Storable a => Stored a -> Stored a -> StoredList a -> IO (StoredList a)-slistReplaceS rm sx prev@(Stored (Ref st _) _) = wrappedStore st $ ListItem [prev] (Just sx) (findSListRef rm prev)--findSListRef :: Stored a -> StoredList a -> Maybe (StoredList a)-findSListRef _ (Stored _ ListNil) = Nothing-findSListRef x cur | listItem (fromStored cur) == Just x = Just cur- | otherwise = listToMaybe $ catMaybes $ map (findSListRef x) $ listPrev $ fromStored cur--{- TODO-mapFromSList :: Storable a => StoredList a -> Map RefDigest (Stored a)-mapFromSList list = helper list M.empty- where helper :: Storable a => StoredList a -> Map RefDigest (Stored a) -> Map RefDigest (Stored a)- helper (Stored _ ListNil) cur = cur- helper (Stored _ (ListItem (Just rref) _ (Just x) rest)) cur =- let rxref = case load rref of- ListItem _ _ (Just rx) _ -> sameType rx x $ storedRef rx- _ -> error "mapFromSList: malformed list"- in helper rest $ case M.lookup (refDigest $ storedRef x) cur of- Nothing -> M.insert (refDigest rxref) x cur- Just x' -> M.insert (refDigest rxref) x' cur- helper (Stored _ (ListItem _ _ _ rest)) cur = helper rest cur- sameType :: a -> a -> b -> b- sameType _ _ x = x--updateOld :: Map RefDigest (Stored a) -> Stored a -> Stored a-updateOld m x = fromMaybe x $ M.lookup (refDigest $ storedRef x) m---data StoreUpdate a = StoreKeep- | StoreReplace a- | StoreRemove--withStoredListItem :: (Storable a) => (a -> Bool) -> StoredList a -> (a -> IO (StoreUpdate a)) -> IO (StoredList a)-withStoredListItem p list f = withStoredListItemS (p . fromStored) list (suMap (wrappedStore $ storedStorage list) <=< f . fromStored)- where suMap :: Monad m => (a -> m b) -> StoreUpdate a -> m (StoreUpdate b)- suMap _ StoreKeep = return StoreKeep- suMap g (StoreReplace x) = return . StoreReplace =<< g x- suMap _ StoreRemove = return StoreRemove--withStoredListItemS :: (Storable a) => (Stored a -> Bool) -> StoredList a -> (Stored a -> IO (StoreUpdate (Stored a))) -> IO (StoredList a)-withStoredListItemS p list f = do- case find p $ storedFromSList list of- Just sx -> f sx >>= \case StoreKeep -> return list- StoreReplace nx -> slistReplaceS sx nx list- StoreRemove -> slistRemove sx list- Nothing -> return list--}
src/Erebos/Storage/Merge.hs view
@@ -97,13 +97,16 @@ doLookup x +-- |Returns list of sets starting with the set of given objects and+-- intcrementally adding parents. generations :: Storable a => [Stored a] -> [Set (Stored a)] generations = unfoldr gen . (,S.empty)- where gen (hs, cur) = case filter (`S.notMember` cur) $ previous =<< hs of+ where gen (hs, cur) = case filter (`S.notMember` cur) hs of [] -> Nothing added -> let next = foldr S.insert cur added- in Just (next, (added, next))+ in Just (next, (previous =<< added, next)) +-- |Returns set containing all given objects and their ancestors ancestors :: Storable a => [Stored a] -> Set (Stored a) ancestors = last . (S.empty:) . generations