erebos 0.1.7 → 0.1.8
raw patch · 15 files changed
+751/−250 lines, 15 filesdep ~basedep ~containersdep ~crypton
Dependency ranges changed: base, containers, crypton, fsnotify, hashable, hashtables, network, template-haskell, time
Files
- CHANGELOG.md +7/−0
- README.md +7/−0
- erebos.cabal +15/−14
- main/Main.hs +65/−15
- main/Test.hs +58/−8
- src/Erebos/Chatroom.hs +38/−25
- src/Erebos/Conversation.hs +5/−0
- src/Erebos/Discovery.hs +215/−83
- src/Erebos/ICE.chs +39/−10
- src/Erebos/ICE/pjproject.c +76/−37
- src/Erebos/ICE/pjproject.h +6/−1
- src/Erebos/Network.hs +52/−11
- src/Erebos/Network/ifaddrs.c +110/−0
- src/Erebos/Service.hs +7/−1
- src/Erebos/Storage.hs +51/−45
CHANGELOG.md view
@@ -1,5 +1,12 @@ # Revision history for erebos +## 0.1.8 -- 2025-03-28++* Discovery service without requiring ICE support+* Added `/delete` command to delete chatrooms for current user+* Ignore record items with unexpected type+* Support GHC 9.12+ ## 0.1.7 -- 2024-10-30 * Chatroom-specific identity
README.md view
@@ -137,6 +137,13 @@ : Leave the chatroom. User will no longer be listed as a member and erebos tool will no longer collect message of this chatroom. +`/delete`+: Delete the chatroom; this action is only synchronized with devices belonging+to the current user and does not affect the chatroom state for others. Due to+the storage design, the chatroom data will not be purged from the local state+history, but the chatroom will no longer be listed as available and no futher+updates for this chatroom will be collected or shared with other peers.+ ### Add contacts To ensure the identity of the contact and prevent man-in-the-middle attack,
erebos.cabal view
@@ -1,20 +1,20 @@ Cabal-Version: 3.0 Name: erebos-Version: 0.1.7+Version: 0.1.8 Synopsis: Decentralized messaging and synchronization Description: Library and simple CLI interface implementing the Erebos identity management, decentralized messaging and synchronization protocol, along with local storage.- .+ Erebos identity is based on locally stored cryptographic keys, all communication is end-to-end encrypted. Multiple devices can be attached to the same identity, after which they function interchangeably, without any one being in any way "primary"; messages and other state data are then synchronized automatically whenever the devices are able to connect with one another.- .+ See README for usage of the CLI tool. License: BSD-3-Clause License-File: LICENSE@@ -40,11 +40,12 @@ source-repository head type: git- location: git://erebosprotocol.net/erebos+ location: https://code.erebosprotocol.net/erebos common common ghc-options: -Wall+ -Wno-x-partial -fdefer-typed-holes if flag(ci)@@ -54,7 +55,7 @@ -Wno-error=unused-imports build-depends:- base ^>= { 4.15, 4.16, 4.17, 4.18, 4.19, 4.20 },+ base ^>= { 4.15, 4.16, 4.17, 4.18, 4.19, 4.20, 4.21 }, default-extensions: DefaultSignatures@@ -98,6 +99,7 @@ Erebos.Chatroom Erebos.Contact Erebos.Conversation+ Erebos.Discovery Erebos.Identity Erebos.Message Erebos.Network@@ -126,7 +128,6 @@ if flag(ice) exposed-modules:- Erebos.Discovery Erebos.ICE c-sources: src/Erebos/ICE/pjproject.c@@ -142,21 +143,21 @@ binary >=0.8 && <0.11, bytestring >=0.10 && <0.13, clock >=0.8 && < 0.9,- containers >= 0.6 && <0.8,- crypton ^>= { 1.0 },+ containers ^>= { 0.6, 0.7, 0.8 },+ crypton ^>= { 0.34, 1.0 }, deepseq >= 1.4 && <1.6, directory >= 1.3 && <1.4, filepath >=1.4 && <1.6,- fsnotify ^>= { 0.4 },- hashable >=1.3 && <1.5,- hashtables >=1.2 && <1.4,+ fsnotify ^>= { 0.3, 0.4 },+ hashable ^>= { 1.3, 1.4, 1.5 },+ hashtables ^>= { 1.2, 1.3, 1.4 }, iproute >=1.7.12 && <1.8, memory >=0.14 && <0.19, mtl >=2.2 && <2.4,- network >= 3.1 && <3.2,+ network ^>= { 3.1, 3.2 }, stm >=2.5 && <2.6, text >= 1.2 && <2.2,- time >= 1.8 && <1.14,+ time ^>= { 1.8, 1.9, 1.10, 1.11, 1.12, 1.13, 1.14 }, uuid >=1.3 && <1.4, zlib >=0.6 && <0.8 @@ -194,7 +195,7 @@ mtl, network, process >=1.6 && <1.7,- template-haskell ^>= { 2.17, 2.18, 2.19, 2.20, 2.21, 2.22 },+ template-haskell ^>= { 2.17, 2.18, 2.19, 2.20, 2.21, 2.22, 2.23 }, text, time, transformers >= 0.5 && <0.7,
main/Main.hs view
@@ -40,8 +40,8 @@ import Erebos.Contact import Erebos.Chatroom import Erebos.Conversation-#ifdef ENABLE_ICE_SUPPORT import Erebos.Discovery+#ifdef ENABLE_ICE_SUPPORT import Erebos.ICE #endif import Erebos.Identity@@ -102,10 +102,8 @@ True "create contacts with network peers" , ServiceOption "dm" (someService @DirectMessage Proxy) True "direct messages"-#ifdef ENABLE_ICE_SUPPORT , ServiceOption "discovery" (someService @DiscoveryService Proxy) True "peer discovery"-#endif ] options :: [OptDescr (Options -> Options)]@@ -125,6 +123,20 @@ , Option [] ["chatroom-auto-subscribe"] (ReqArg (\count -> \opts -> opts { optChatroomAutoSubscribe = Just (read count) }) "<count>") "automatically subscribe for up to <count> chatrooms"+#ifdef ENABLE_ICE_SUPPORT+ , Option [] [ "discovery-stun-port" ]+ (ReqArg (\value -> serviceAttr $ \attrs -> attrs { discoveryStunPort = Just (read value) }) "<port>")+ "offer specified <port> to discovery peers for STUN protocol"+ , Option [] [ "discovery-stun-server" ]+ (ReqArg (\value -> serviceAttr $ \attrs -> attrs { discoveryStunServer = Just (read value) }) "<server>")+ "offer <server> (domain name or IP address) to discovery peers for STUN protocol"+ , Option [] [ "discovery-turn-port" ]+ (ReqArg (\value -> serviceAttr $ \attrs -> attrs { discoveryTurnPort = Just (read value) }) "<port>")+ "offer specified <port> to discovery peers for TURN protocol"+ , Option [] [ "discovery-turn-server" ]+ (ReqArg (\value -> serviceAttr $ \attrs -> attrs { discoveryTurnServer = Just (read value) }) "<server>")+ "offer <server> (domain name or IP address) to discovery peers for TURN protocol"+#endif , Option [] ["dm-bot-echo"] (ReqArg (\prefix -> \opts -> opts { optDmBotEcho = Just (T.pack prefix) }) "<prefix>") "automatically reply to direct messages with the same text prefixed with <prefix>"@@ -135,8 +147,17 @@ (NoArg $ \opts -> opts { optShowVersion = True }) "show version and exit" ]- where so f opts = opts { optServer = f $ optServer opts }+ where+ so f opts = opts { optServer = f $ optServer opts } + updateService :: Service s => (ServiceAttributes s -> ServiceAttributes s) -> SomeService -> SomeService+ updateService f some@(SomeService proxy attrs)+ | Just f' <- cast f = SomeService proxy (f' attrs)+ | otherwise = some++ serviceAttr :: Service s => (ServiceAttributes s -> ServiceAttributes s) -> Options -> Options+ serviceAttr f opts = opts { optServices = map (\sopt -> sopt { soptService = updateService f (soptService sopt) }) (optServices opts) }+ servicesOptions :: [OptDescr (Options -> Options)] servicesOptions = concatMap helper $ "all" : map soptName availableServices where@@ -480,6 +501,7 @@ , ("peer-add-public", cmdPeerAddPublic) , ("peer-drop", cmdPeerDrop) , ("send", cmdSend)+ , ("delete", cmdDelete) , ("update-identity", cmdUpdateIdentity) , ("attach", cmdAttach) , ("attach-accept", cmdAttachAccept)@@ -492,9 +514,9 @@ , ("contact-reject", cmdContactReject) , ("conversations", cmdConversations) , ("details", cmdDetails)-#ifdef ENABLE_ICE_SUPPORT , ("discovery-init", cmdDiscoveryInit) , ("discovery", cmdDiscovery)+#ifdef ENABLE_ICE_SUPPORT , ("ice-create", cmdIceCreate) , ("ice-destroy", cmdIceDestroy) , ("ice-show", cmdIceShow)@@ -611,6 +633,11 @@ liftIO $ putStrLn $ formatMessage tzone msg Nothing -> return () +cmdDelete :: Command+cmdDelete = void $ do+ deleteConversation =<< getSelectedConversation+ modify $ \s -> s { csContext = NoContext }+ cmdHistory :: Command cmdHistory = void $ do conv <- getSelectedConversation@@ -657,7 +684,7 @@ watchChatrooms h $ \set -> \case Nothing -> do- let chatroomList = fromSetBy (comparing roomStateData) set+ let chatroomList = filter (not . roomStateDeleted) $ fromSetBy (comparing roomStateData) set (subscribed, notSubscribed) = partition roomStateSubscribe chatroomList subscribedNum = length subscribed @@ -717,7 +744,7 @@ cmdChatrooms = do ensureWatchedChatrooms chatroomSetVar <- asks ciChatroomSetVar- chatroomList <- fromSetBy (comparing roomStateData) <$> liftIO (readMVar chatroomSetVar)+ chatroomList <- filter (not . roomStateDeleted) . fromSetBy (comparing roomStateData) <$> liftIO (readMVar chatroomSetVar) set <- asks ciSetContextOptions set $ map SelectedChatroom chatroomList forM_ (zip [1..] chatroomList) $ \(i :: Int, rstate) -> do@@ -838,8 +865,6 @@ , map (BC.unpack . showRefDigest . refDigest . storedRef) $ idExtDataF cpid ] -#ifdef ENABLE_ICE_SUPPORT- cmdDiscoveryInit :: Command cmdDiscoveryInit = void $ do server <- asks ciServer@@ -850,7 +875,7 @@ [] -> ("discovery.erebosprotocol.net", show discoveryPort) addr:_ <- liftIO $ getAddrInfo (Just $ defaultHints { addrSocketType = Datagram }) (Just hostname) (Just port) peer <- liftIO $ serverPeer server (addrAddress addr)- sendToPeer peer $ DiscoverySelf (T.pack "ICE") 0+ sendToPeer peer $ DiscoverySelf [ T.pack "ICE" ] Nothing modify $ \s -> s { csIcePeer = Just peer } cmdDiscovery :: Command@@ -867,14 +892,39 @@ Right _ -> return () Left err -> eprint err +#ifdef ENABLE_ICE_SUPPORT+ cmdIceCreate :: Command cmdIceCreate = do- role <- asks ciLine >>= return . \case- 'm':_ -> PjIceSessRoleControlling- 's':_ -> PjIceSessRoleControlled- _ -> PjIceSessRoleUnknown+ let getRole = \case+ 'm':_ -> PjIceSessRoleControlling+ 's':_ -> PjIceSessRoleControlled+ _ -> PjIceSessRoleUnknown++ ( role, stun, turn ) <- asks (words . ciLine) >>= \case+ [] -> return ( PjIceSessRoleControlling, Nothing, Nothing )+ [ role ] -> return+ ( getRole role, Nothing, Nothing )+ [ role, server ] -> return+ ( getRole role+ , Just ( T.pack server, 0 )+ , Just ( T.pack server, 0 )+ )+ [ role, server, port ] -> return+ ( getRole role+ , Just ( T.pack server, read port )+ , Just ( T.pack server, read port )+ )+ [ role, stunServer, stunPort, turnServer, turnPort ] -> return+ ( getRole role+ , Just ( T.pack stunServer, read stunPort )+ , Just ( T.pack turnServer, read turnPort )+ )+ _ -> throwError "invalid parameters"+ eprint <- asks ciPrint- sess <- liftIO $ iceCreate role $ eprint <=< iceShow+ Just cfg <- liftIO $ iceCreateConfig stun turn+ sess <- liftIO $ iceCreateSession cfg role $ eprint <=< iceShow modify $ \s -> s { csIceSessions = sess : csIceSessions s } cmdIceDestroy :: Command
main/Test.hs view
@@ -1,3 +1,5 @@+{-# LANGUAGE OverloadedStrings #-}+ module Test ( runTestTool, ) where@@ -34,6 +36,7 @@ import Erebos.Attach import Erebos.Chatroom import Erebos.Contact+import Erebos.Discovery import Erebos.Identity import Erebos.Message import Erebos.Network@@ -255,6 +258,7 @@ , ("head-watch", cmdHeadWatch) , ("head-unwatch", cmdHeadUnwatch) , ("create-identity", cmdCreateIdentity)+ , ("identity-info", cmdIdentityInfo) , ("start-server", cmdStartServer) , ("stop-server", cmdStopServer) , ("peer-add", cmdPeerAdd)@@ -283,6 +287,7 @@ , ("dm-list-peer", cmdDmListPeer) , ("dm-list-contact", cmdDmListContact) , ("chatroom-create", cmdChatroomCreate)+ , ("chatroom-delete", cmdChatroomDelete) , ("chatroom-list-local", cmdChatroomListLocal) , ("chatroom-watch-local", cmdChatroomWatchLocal) , ("chatroom-set-name", cmdChatroomSetName)@@ -293,6 +298,7 @@ , ("chatroom-join-as", cmdChatroomJoinAs) , ("chatroom-leave", cmdChatroomLeave) , ("chatroom-message-send", cmdChatroomMessageSend)+ , ("discovery-connect", cmdDiscoveryConnect) ] cmdStore :: Command@@ -443,27 +449,53 @@ , lsOther = [] } initTestHead h+ cmdOut $ unwords [ "create-identity-done", "ref", show $ refDigest $ storedRef $ lsIdentity $ headObject h ] +cmdIdentityInfo :: Command+cmdIdentityInfo = do+ st <- asks tiStorage+ [ tref ] <- asks tiParams+ Just ref <- liftIO $ readRef st $ encodeUtf8 tref+ let sidata = wrappedLoad ref+ idata = fromSigned sidata+ cmdOut $ unwords $ concat+ [ [ "identity-info" ]+ , [ "ref", T.unpack tref ]+ , [ "base", show $ refDigest $ storedRef $ eiddStoredBase sidata ]+ , maybe [] (\owner -> [ "owner", show $ refDigest $ storedRef owner ]) $ eiddOwner idata+ , maybe [] (\name -> [ "name", T.unpack name ]) $ eiddName idata+ ]+ cmdStartServer :: Command cmdStartServer = do out <- asks tiOutput + let parseParams = \case+ (name : value : rest)+ | name == "services" -> T.splitOn "," value+ | otherwise -> parseParams rest+ _ -> []+ serviceNames <- parseParams <$> asks tiParams+ h <- getOrLoadHead rsPeers <- liftIO $ newMVar (1, [])- 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- , someService @SyncService Proxy- , someService @ChatroomService Proxy- , someServiceAttr $ (defaultServiceAttributes Proxy)+ services <- forM serviceNames $ \case+ "attach" -> return $ someServiceAttr $ pairingAttributes (Proxy @AttachService) out rsPeers "attach"+ "chatroom" -> return $ someService @ChatroomService Proxy+ "contact" -> return $ someServiceAttr $ pairingAttributes (Proxy @ContactService) out rsPeers "contact"+ "discovery" -> return $ someService @DiscoveryService Proxy+ "dm" -> return $ someServiceAttr $ directMessageAttributes out+ "sync" -> return $ someService @SyncService Proxy+ "test" -> return $ someServiceAttr $ (defaultServiceAttributes Proxy) { testMessageReceived = \obj otype len sref -> do liftIO $ do void $ store (headStorage h) obj outLine out $ unwords ["test-message-received", otype, len, sref] }- ]+ sname -> throwError $ "unknown service `" <> T.unpack sname <> "'" + rsServer <- liftIO $ startServer defaultServerOptions h (B.hPutStr stderr . (`BC.snoc` '\n') . BC.pack) services+ rsPeerThread <- liftIO $ forkIO $ void $ forever $ do peer <- getNextPeerChange rsServer @@ -733,6 +765,13 @@ room <- createChatroom (Just name) Nothing cmdOut $ unwords $ "chatroom-create-done" : chatroomInfo room +cmdChatroomDelete :: Command+cmdChatroomDelete = do+ [ cid ] <- asks tiParams+ sdata <- getChatroomStateData cid+ deleteChatroomByStateData sdata+ cmdOut $ unwords [ "chatroom-delete-done", T.unpack cid ]+ getChatroomStateData :: Text -> CommandM (Stored ChatroomStateData) getChatroomStateData tref = do st <- asks tiStorage@@ -835,3 +874,14 @@ [cid, msg] <- asks tiParams to <- getChatroomStateData cid void $ sendChatroomMessageByStateData to msg++cmdDiscoveryConnect :: Command+cmdDiscoveryConnect = do+ st <- asks tiStorage+ [ tref ] <- asks tiParams+ Just ref <- liftIO $ readRef st $ encodeUtf8 tref++ Just RunningServer {..} <- gets tsServer+ peers <- liftIO $ getCurrentPeerList rsServer+ forM_ peers $ \peer -> do+ sendToPeer peer $ DiscoverySearch ref
src/Erebos/Chatroom.hs view
@@ -6,6 +6,7 @@ ChatroomState(..), ChatroomStateData(..), createChatroom,+ deleteChatroomByStateData, updateChatroomByStateData, listChatrooms, findChatroomByRoomData,@@ -206,9 +207,8 @@ else [] mdata <- mstore =<< sign secret =<< mstore ChatMessageData {..}- mergeSorted . (:[]) <$> mstore ChatroomStateData+ mergeSorted . (:[]) <$> mstore emptyChatroomStateData { rsdPrev = roomStateData cstate- , rsdRoom = [] , rsdSubscribe = Just (not mdLeave) , rsdIdentity = mbIdentity , rsdMessages = [ mdata ]@@ -218,15 +218,27 @@ data ChatroomStateData = ChatroomStateData { rsdPrev :: [Stored ChatroomStateData] , rsdRoom :: [Stored (Signed ChatroomData)]+ , rsdDelete :: Bool , rsdSubscribe :: Maybe Bool , rsdIdentity :: Maybe UnifiedIdentity , rsdMessages :: [Stored (Signed ChatMessageData)] } +emptyChatroomStateData :: ChatroomStateData+emptyChatroomStateData = ChatroomStateData+ { rsdPrev = []+ , rsdRoom = []+ , rsdDelete = False+ , rsdSubscribe = Nothing+ , rsdIdentity = Nothing+ , rsdMessages = []+ }+ data ChatroomState = ChatroomState { roomStateData :: [Stored ChatroomStateData] , roomStateRoom :: Maybe Chatroom , roomStateMessageData :: [Stored (Signed ChatMessageData)]+ , roomStateDeleted :: Bool , roomStateSubscribe :: Bool , roomStateIdentity :: Maybe UnifiedIdentity , roomStateMessages :: [ChatMessage]@@ -236,6 +248,7 @@ store' ChatroomStateData {..} = storeRec $ do forM_ rsdPrev $ storeRef "PREV" forM_ rsdRoom $ storeRef "room"+ when rsdDelete $ storeEmpty "delete" forM_ rsdSubscribe $ storeInt "subscribe" . bool @Int 0 1 forM_ rsdIdentity $ storeRef "id" . idExtData forM_ rsdMessages $ storeRef "msg"@@ -243,6 +256,7 @@ load' = loadRec $ do rsdPrev <- loadRefs "PREV" rsdRoom <- loadRefs "room"+ rsdDelete <- isJust <$> loadMbEmpty "delete" rsdSubscribe <- fmap ((/=) @Int 0) <$> loadMbInt "subscribe" rsdIdentity <- loadMbUnifiedIdentity "id" rsdMessages <- loadRefs "msg"@@ -257,7 +271,8 @@ roomStateMessageData = filterAncestors $ concat $ flip findProperty roomStateData $ \case ChatroomStateData {..} | null rsdMessages -> Nothing | otherwise -> Just rsdMessages- roomStateSubscribe = fromMaybe False $ findPropertyFirst rsdSubscribe roomStateData+ roomStateDeleted = any (rsdDelete . fromStored) roomStateData+ roomStateSubscribe = not roomStateDeleted && (fromMaybe False $ findPropertyFirst rsdSubscribe roomStateData) roomStateIdentity = findPropertyFirst rsdIdentity roomStateData roomStateMessages = threadToListSince [] $ concatMap (rsdMessages . fromStored) roomStateData in ChatroomState {..}@@ -272,12 +287,9 @@ (secret, rdKey) <- liftIO . generateKeys =<< getStorage let rdPrev = [] rdata <- mstore =<< sign secret =<< mstore ChatroomData {..}- cstate <- mergeSorted . (:[]) <$> mstore ChatroomStateData- { rsdPrev = []- , rsdRoom = [ rdata ]+ cstate <- mergeSorted . (:[]) <$> mstore emptyChatroomStateData+ { rsdRoom = [ rdata ] , rsdSubscribe = Just True- , rsdIdentity = Nothing- , rsdMessages = [] } updateLocalHead $ updateSharedState $ \rooms -> do@@ -303,6 +315,17 @@ return (roomSet, Just upd) [] -> return (roomSet, Nothing) +deleteChatroomByStateData+ :: (MonadStorage m, MonadHead LocalState m, MonadError String m)+ => Stored ChatroomStateData -> m ()+deleteChatroomByStateData lookupData = void $ findAndUpdateChatroomState $ \cstate -> do+ guard $ any (lookupData `precedesOrEquals`) $ roomStateData cstate+ Just $ do+ mergeSorted . (:[]) <$> mstore emptyChatroomStateData+ { rsdPrev = roomStateData cstate+ , rsdDelete = True+ }+ updateChatroomByStateData :: (MonadStorage m, MonadHead LocalState m, MonadError String m) => Stored ChatroomStateData@@ -320,17 +343,16 @@ , rdDescription = newDesc , rdKey = roomKey room }- mergeSorted . (:[]) <$> mstore ChatroomStateData+ mergeSorted . (:[]) <$> mstore emptyChatroomStateData { rsdPrev = roomStateData cstate , rsdRoom = [ rdata ] , rsdSubscribe = Just True- , rsdIdentity = Nothing- , rsdMessages = [] } listChatrooms :: MonadHead LocalState m => m [ChatroomState]-listChatrooms = fromSetBy (comparing $ roomName <=< roomStateRoom) .+listChatrooms = filter (not . roomStateDeleted) .+ fromSetBy (comparing $ roomName <=< roomStateRoom) . lookupSharedValue . lsShared . fromStored <$> getLocalHead findChatroom :: MonadHead LocalState m => (ChatroomState -> Bool) -> m (Maybe ChatroomState)@@ -351,12 +373,9 @@ chatroomSetSubscribe lookupData subscribe = void $ findAndUpdateChatroomState $ \cstate -> do guard $ any (lookupData `precedesOrEquals`) $ roomStateData cstate Just $ do- mergeSorted . (:[]) <$> mstore ChatroomStateData+ mergeSorted . (:[]) <$> mstore emptyChatroomStateData { rsdPrev = roomStateData cstate- , rsdRoom = [] , rsdSubscribe = Just subscribe- , rsdIdentity = Nothing- , rsdMessages = [] } chatroomMembers :: ChatroomState -> [ ComposedIdentity ]@@ -419,7 +438,7 @@ return $ makeChatroomDiff lastList curList chatroomSetToList :: Set ChatroomState -> [(Stored ChatroomStateData, ChatroomState)]-chatroomSetToList = map (cmp &&& id) . fromSetBy (comparing cmp)+chatroomSetToList = map (cmp &&& id) . filter (not . roomStateDeleted) . fromSetBy (comparing cmp) where cmp :: ChatroomState -> Stored ChatroomStateData cmp = head . filterAncestors . concatMap storedRoots . toComponents@@ -517,12 +536,9 @@ -- update local state only if we got roomInfo not present there if roomInfo `notElem` prevRoom && roomInfo `elem` room then do- sdata <- mstore ChatroomStateData+ sdata <- mstore emptyChatroomStateData { rsdPrev = prev , rsdRoom = room- , rsdSubscribe = Nothing- , rsdIdentity = Nothing- , rsdMessages = [] } storeSetAddComponent sdata set else return set@@ -562,11 +578,8 @@ -- update local state only if subscribed and we got some new messages if roomStateSubscribe prev && messages /= prevMessages then do- sdata <- mstore ChatroomStateData+ sdata <- mstore emptyChatroomStateData { rsdPrev = prevData- , rsdRoom = []- , rsdSubscribe = Nothing- , rsdIdentity = Nothing , rsdMessages = messages } storeSetAddComponent sdata set
src/Erebos/Conversation.hs view
@@ -18,6 +18,7 @@ conversationHistory, sendMessage,+ deleteConversation, ) where import Control.Monad.Except@@ -103,3 +104,7 @@ sendMessage :: (MonadHead LocalState m, MonadError String m) => Conversation -> Text -> m (Maybe Message) sendMessage (DirectMessageConversation thread) text = fmap Just $ DirectMessageMessage <$> (fromStored <$> sendDirectMessage (msgPeer thread) text) <*> pure False sendMessage (ChatroomConversation rstate) text = sendChatroomMessage rstate text >> return Nothing++deleteConversation :: (MonadHead LocalState m, MonadError String m) => Conversation -> m ()+deleteConversation (DirectMessageConversation _) = throwError "deleting direct message conversation is not supported"+deleteConversation (ChatroomConversation rstate) = deleteChatroomByStateData (head $ roomStateData rstate)
src/Erebos/Discovery.hs view
@@ -1,5 +1,8 @@+{-# LANGUAGE CPP #-}+ module Erebos.Discovery ( DiscoveryService(..),+ DiscoveryAttributes(..), DiscoveryConnection(..) ) where @@ -8,54 +11,82 @@ import Control.Monad.Except import Control.Monad.Reader +import Data.IP qualified as IP import Data.Map.Strict (Map)-import qualified Data.Map.Strict as M+import Data.Map.Strict qualified as M import Data.Maybe import Data.Text (Text)-import qualified Data.Text as T+import Data.Text qualified as T+import Data.Word import Network.Socket +#ifdef ENABLE_ICE_SUPPORT import Erebos.ICE+#endif import Erebos.Identity import Erebos.Network import Erebos.Service import Erebos.Storage -keepaliveSeconds :: Int-keepaliveSeconds = 20+data DiscoveryService+ = DiscoverySelf [ Text ] (Maybe Int)+ | DiscoveryAcknowledged [ Text ] (Maybe Text) (Maybe Word16) (Maybe Text) (Maybe Word16)+ | DiscoverySearch Ref+ | DiscoveryResult Ref [ Text ]+ | DiscoveryConnectionRequest DiscoveryConnection+ | DiscoveryConnectionResponse DiscoveryConnection +data DiscoveryAttributes = DiscoveryAttributes+ { discoveryStunPort :: Maybe Word16+ , discoveryStunServer :: Maybe Text+ , discoveryTurnPort :: Maybe Word16+ , discoveryTurnServer :: Maybe Text+ } -data DiscoveryService = DiscoverySelf Text Int- | DiscoveryAcknowledged Text- | DiscoverySearch Ref- | DiscoveryResult Ref (Maybe Text)- | DiscoveryConnectionRequest DiscoveryConnection- | DiscoveryConnectionResponse DiscoveryConnection+defaultDiscoveryAttributes :: DiscoveryAttributes+defaultDiscoveryAttributes = DiscoveryAttributes+ { discoveryStunPort = Nothing+ , discoveryStunServer = Nothing+ , discoveryTurnPort = Nothing+ , discoveryTurnServer = Nothing+ } data DiscoveryConnection = DiscoveryConnection { dconnSource :: Ref , dconnTarget :: Ref , dconnAddress :: Maybe Text- , dconnIceSession :: Maybe IceRemoteInfo+#ifdef ENABLE_ICE_SUPPORT+ , dconnIceInfo :: Maybe IceRemoteInfo+#else+ , dconnIceInfo :: Maybe (Stored Object)+#endif } emptyConnection :: Ref -> Ref -> DiscoveryConnection-emptyConnection source target = DiscoveryConnection source target Nothing Nothing+emptyConnection dconnSource dconnTarget = DiscoveryConnection {..}+ where+ dconnAddress = Nothing+ dconnIceInfo = Nothing instance Storable DiscoveryService where store' x = storeRec $ do case x of- DiscoverySelf addr priority -> do- storeText "self" addr- storeInt "priority" priority- DiscoveryAcknowledged addr -> do- storeText "ack" addr+ DiscoverySelf addrs priority -> do+ mapM_ (storeText "self") addrs+ mapM_ (storeInt "priority") priority+ DiscoveryAcknowledged addrs stunServer stunPort turnServer turnPort -> do+ if null addrs then storeEmpty "ack"+ else mapM_ (storeText "ack") addrs+ storeMbText "stun-server" stunServer+ storeMbInt "stun-port" stunPort+ storeMbText "turn-server" turnServer+ storeMbInt "turn-port" turnPort DiscoverySearch ref -> storeRawRef "search" ref DiscoveryResult ref addr -> do storeRawRef "result" ref- storeMbText "address" addr+ mapM_ (storeText "address") addr DiscoveryConnectionRequest conn -> storeConnection "request" conn DiscoveryConnectionResponse conn -> storeConnection "response" conn @@ -64,18 +95,28 @@ storeRawRef "source" $ dconnSource conn storeRawRef "target" $ dconnTarget conn storeMbText "address" $ dconnAddress conn- storeMbRef "ice-session" $ dconnIceSession conn+ storeMbRef "ice-info" $ dconnIceInfo conn load' = loadRec $ msum- [ DiscoverySelf- <$> loadText "self"- <*> loadInt "priority"- , DiscoveryAcknowledged- <$> loadText "ack"+ [ do+ addrs <- loadTexts "self"+ guard (not $ null addrs)+ DiscoverySelf addrs+ <$> loadMbInt "priority"+ , do+ addrs <- loadTexts "ack"+ mbEmpty <- loadMbEmpty "ack"+ guard (not (null addrs) || isJust mbEmpty)+ DiscoveryAcknowledged+ <$> pure addrs+ <*> loadMbText "stun-server"+ <*> loadMbInt "stun-port"+ <*> loadMbText "turn-server"+ <*> loadMbInt "turn-port" , DiscoverySearch <$> loadRawRef "search" , DiscoveryResult <$> loadRawRef "result"- <*> loadMbText "address"+ <*> loadTexts "address" , loadConnection "request" DiscoveryConnectionRequest , loadConnection "response" DiscoveryConnectionResponse ]@@ -86,113 +127,183 @@ <$> loadRawRef "source" <*> loadRawRef "target" <*> loadMbText "address"- <*> loadMbRef "ice-session"+ <*> loadMbRef "ice-info" data DiscoveryPeer = DiscoveryPeer { dpPriority :: Int , dpPeer :: Maybe Peer- , dpAddress :: Maybe Text+ , dpAddress :: [ Text ]+#ifdef ENABLE_ICE_SUPPORT , dpIceSession :: Maybe IceSession+#endif } instance Service DiscoveryService where- serviceID _ = mkServiceID "dd59c89c-69cc-4703-b75b-4ddcd4b3c23b"+ serviceID _ = mkServiceID "dd59c89c-69cc-4703-b75b-4ddcd4b3c23c" + type ServiceAttributes DiscoveryService = DiscoveryAttributes+ defaultServiceAttributes _ = defaultDiscoveryAttributes++#ifdef ENABLE_ICE_SUPPORT+ type ServiceState DiscoveryService = Maybe IceConfig+ emptyServiceState _ = Nothing+#endif+ type ServiceGlobalState DiscoveryService = Map RefDigest DiscoveryPeer emptyServiceGlobalState _ = M.empty serviceHandler msg = case fromStored msg of- DiscoverySelf addr priority -> do+ DiscoverySelf addrs priority -> do pid <- asks svcPeerIdentity peer <- asks svcPeer let insertHelper new old | dpPriority new > dpPriority old = new | otherwise = old- mbaddr <- case words (T.unpack addr) of- [ipaddr, port] | DatagramAddress paddr <- peerAddress peer -> do+ matchedAddrs <- fmap catMaybes $ forM addrs $ \addr -> if+ | addr == T.pack "ICE" -> do+ return $ Just addr++ | [ ipaddr, port ] <- words (T.unpack addr)+ , DatagramAddress paddr <- peerAddress peer -> do saddr <- liftIO $ head <$> getAddrInfo (Just $ defaultHints { addrSocketType = Datagram }) (Just ipaddr) (Just port) return $ if paddr == addrAddress saddr then Just addr else Nothing- _ -> return Nothing++ | otherwise -> return Nothing+ forM_ (idDataF =<< unfoldOwners pid) $ \s ->- svcModifyGlobal $ M.insertWith insertHelper (refDigest $ storedRef s) $- DiscoveryPeer priority (Just peer) mbaddr Nothing- replyPacket $ DiscoveryAcknowledged $ fromMaybe (T.pack "ICE") mbaddr+ svcModifyGlobal $ M.insertWith insertHelper (refDigest $ storedRef s) DiscoveryPeer+ { dpPriority = fromMaybe 0 priority+ , dpPeer = Just peer+ , dpAddress = addrs+#ifdef ENABLE_ICE_SUPPORT+ , dpIceSession = Nothing+#endif+ }+ attrs <- asks svcAttributes+ replyPacket $ DiscoveryAcknowledged matchedAddrs+ (discoveryStunServer attrs)+ (discoveryStunPort attrs)+ (discoveryTurnServer attrs)+ (discoveryTurnPort attrs) - DiscoveryAcknowledged addr -> do- when (addr == T.pack "ICE") $ do- -- keep-alive packet from behind NAT- peer <- asks svcPeer- liftIO $ void $ forkIO $ do- threadDelay (keepaliveSeconds * 1000 * 1000)- res <- runExceptT $ sendToPeer peer $ DiscoverySelf addr 0- case res of- Right _ -> return ()- Left err -> putStrLn $ "Discovery: failed to send keep-alive: " ++ err+ DiscoveryAcknowledged _ stunServer stunPort turnServer turnPort -> do+#ifdef ENABLE_ICE_SUPPORT+ paddr <- asks (peerAddress . svcPeer) >>= return . \case+ (DatagramAddress saddr) -> case IP.fromSockAddr saddr of+ Just (IP.IPv6 ipv6, _)+ | (0, 0, 0xffff, ipv4) <- IP.fromIPv6w ipv6+ -> Just $ T.pack $ show (IP.toIPv4w ipv4)+ Just (addr, _)+ -> Just $ T.pack $ show addr+ _ -> Nothing+ _ -> Nothing + let toIceServer Nothing Nothing = Nothing+ toIceServer Nothing (Just port) = ( , port) <$> paddr+ toIceServer (Just server) Nothing = Just ( server, 0 )+ toIceServer (Just server) (Just port) = Just ( server, port )++ cfg <- liftIO $ iceCreateConfig+ (toIceServer stunServer stunPort)+ (toIceServer turnServer turnPort)+ svcSet cfg+#endif+ return ()+ DiscoverySearch ref -> do- addr <- M.lookup (refDigest ref) <$> svcGetGlobal- replyPacket $ DiscoveryResult ref $ fromMaybe (T.pack "ICE") . dpAddress <$> addr+ dpeer <- M.lookup (refDigest ref) <$> svcGetGlobal+ replyPacket $ DiscoveryResult ref $ maybe [] dpAddress dpeer - DiscoveryResult ref Nothing -> do+ DiscoveryResult ref [] -> do svcPrint $ "Discovery: " ++ show (refDigest ref) ++ " not found" - DiscoveryResult ref (Just addr) -> do+ DiscoveryResult ref addrs -> do -- TODO: check if we really requested that server <- asks svcServer- if addr == T.pack "ICE"- then do- self <- svcSelf- peer <- asks svcPeer- ice <- liftIO $ iceCreate PjIceSessRoleControlling $ \ice -> do+ self <- svcSelf+ mbIceConfig <- svcGet+ discoveryPeer <- asks svcPeer+ let runAsService = runPeerService @DiscoveryService discoveryPeer++ liftIO $ void $ forkIO $ forM_ addrs $ \addr -> if+ | addr == T.pack "ICE"+#ifdef ENABLE_ICE_SUPPORT+ , Just config <- mbIceConfig+ -> do+ ice <- iceCreateSession config PjIceSessRoleControlling $ \ice -> do rinfo <- iceRemoteInfo ice- res <- runExceptT $ sendToPeer peer $- DiscoveryConnectionRequest (emptyConnection (storedRef $ idData self) ref) { dconnIceSession = Just rinfo }+ res <- runExceptT $ sendToPeer discoveryPeer $+ DiscoveryConnectionRequest (emptyConnection (storedRef $ idData self) ref) { dconnIceInfo = Just rinfo } case res of Right _ -> return () Left err -> putStrLn $ "Discovery: failed to send connection request: " ++ err - svcModifyGlobal $ M.insert (refDigest ref) $- DiscoveryPeer 0 Nothing Nothing (Just ice)- else do- case words (T.unpack addr) of- [ipaddr, port] -> do- saddr <- liftIO $ head <$>- getAddrInfo (Just $ defaultHints { addrSocketType = Datagram }) (Just ipaddr) (Just port)- peer <- liftIO $ serverPeer server (addrAddress saddr)- svcModifyGlobal $ M.insert (refDigest ref) $- DiscoveryPeer 0 (Just peer) Nothing Nothing+ runAsService $ do+ svcModifyGlobal $ M.insert (refDigest ref) DiscoveryPeer+ { dpPriority = 0+ , dpPeer = Nothing+ , dpAddress = []+ , dpIceSession = Just ice+ }+#else+ -> do+ return ()+#endif - _ -> svcPrint $ "Discovery: invalid address in result: " ++ T.unpack addr+ | [ ipaddr, port ] <- words (T.unpack addr) -> do+ saddr <- head <$>+ getAddrInfo (Just $ defaultHints { addrSocketType = Datagram }) (Just ipaddr) (Just port)+ peer <- serverPeer server (addrAddress saddr)+ runAsService $ do+ svcModifyGlobal $ M.insert (refDigest ref) DiscoveryPeer+ { dpPriority = 0+ , dpPeer = Just peer+ , dpAddress = []+#ifdef ENABLE_ICE_SUPPORT+ , dpIceSession = Nothing+#endif+ } + | otherwise -> do+ runAsService $ do+ svcPrint $ "Discovery: invalid address in result: " ++ T.unpack addr+ DiscoveryConnectionRequest conn -> do self <- svcSelf let rconn = emptyConnection (dconnSource conn) (dconnTarget conn) if refDigest (dconnTarget conn) `elem` (map (refDigest . storedRef) $ idDataF =<< unfoldOwners self) then do+#ifdef ENABLE_ICE_SUPPORT -- request for us, create ICE sesssion server <- asks svcServer peer <- asks svcPeer- liftIO $ void $ iceCreate PjIceSessRoleControlled $ \ice -> do- rinfo <- iceRemoteInfo ice- res <- runExceptT $ sendToPeer peer $ DiscoveryConnectionResponse rconn { dconnIceSession = Just rinfo }- case res of- Right _ -> do- case dconnIceSession conn of- Just prinfo -> iceConnect ice prinfo $ void $ serverPeerIce server ice- Nothing -> putStrLn $ "Discovery: connection request without ICE remote info"- Left err -> putStrLn $ "Discovery: failed to send connection response: " ++ err+ svcGet >>= \case+ Just config -> do+ liftIO $ void $ iceCreateSession config PjIceSessRoleControlled $ \ice -> do+ rinfo <- iceRemoteInfo ice+ res <- runExceptT $ sendToPeer peer $ DiscoveryConnectionResponse rconn { dconnIceInfo = Just rinfo }+ case res of+ Right _ -> do+ case dconnIceInfo conn of+ Just prinfo -> iceConnect ice prinfo $ void $ serverPeerIce server ice+ Nothing -> putStrLn $ "Discovery: connection request without ICE remote info"+ Left err -> putStrLn $ "Discovery: failed to send connection response: " ++ err+ Nothing -> do+ svcPrint $ "Discovery: ICE request from peer without ICE configuration"+#else+ return ()+#endif else do -- request to some of our peers, relay mbdp <- M.lookup (refDigest $ dconnTarget conn) <$> svcGetGlobal case mbdp of Nothing -> replyPacket $ DiscoveryConnectionResponse rconn- Just dp | Just addr <- dpAddress dp -> do- replyPacket $ DiscoveryConnectionResponse rconn { dconnAddress = Just addr }- | Just dpeer <- dpPeer dp -> do- sendToPeer dpeer $ DiscoveryConnectionRequest conn- | otherwise -> svcPrint $ "Discovery: failed to relay connection request"+ Just dp+ | Just dpeer <- dpPeer dp -> do+ sendToPeer dpeer $ DiscoveryConnectionRequest conn+ | otherwise -> svcPrint $ "Discovery: failed to relay connection request" DiscoveryConnectionResponse conn -> do self <- svcSelf@@ -200,6 +311,7 @@ if refDigest (dconnSource conn) `elem` (map (refDigest . storedRef) $ idDataF =<< unfoldOwners self) then do -- response to our request, try to connect to the peer+#ifdef ENABLE_ICE_SUPPORT server <- asks svcServer if | Just addr <- dconnAddress conn , [ipaddr, port] <- words (T.unpack addr) -> do@@ -207,17 +319,37 @@ getAddrInfo (Just $ defaultHints { addrSocketType = Datagram }) (Just ipaddr) (Just port) peer <- liftIO $ serverPeer server (addrAddress saddr) svcModifyGlobal $ M.insert (refDigest $ dconnTarget conn) $- DiscoveryPeer 0 (Just peer) Nothing Nothing+ DiscoveryPeer 0 (Just peer) [] Nothing | Just dp <- M.lookup (refDigest $ dconnTarget conn) dpeers , Just ice <- dpIceSession dp- , Just rinfo <- dconnIceSession conn -> do+ , Just rinfo <- dconnIceInfo conn -> do liftIO $ iceConnect ice rinfo $ void $ serverPeerIce server ice | otherwise -> svcPrint $ "Discovery: connection request failed"+#else+ return ()+#endif else do -- response to relayed request case M.lookup (refDigest $ dconnSource conn) dpeers of Just dp | Just dpeer <- dpPeer dp -> do sendToPeer dpeer $ DiscoveryConnectionResponse conn _ -> svcPrint $ "Discovery: failed to relay connection response"++ serviceNewPeer = do+ server <- asks svcServer+ peer <- asks svcPeer++ let addrToText saddr = do+ ( addr, port ) <- IP.fromSockAddr saddr+ Just $ T.pack $ show addr <> " " <> show port+ addrs <- concat <$> sequence+ [ catMaybes . map addrToText <$> liftIO (getServerAddresses server)+#ifdef ENABLE_ICE_SUPPORT+ , return [ T.pack "ICE" ]+#endif+ ]++ when (not $ null addrs) $ do+ sendToPeer peer $ DiscoverySelf addrs Nothing
src/Erebos/ICE.chs view
@@ -4,9 +4,11 @@ module Erebos.ICE ( IceSession, IceSessionRole(..),+ IceConfig, IceRemoteInfo, - iceCreate,+ iceCreateConfig,+ iceCreateSession, iceDestroy, iceRemoteInfo, iceShow,@@ -17,23 +19,25 @@ ) where import Control.Arrow-import Control.Concurrent.MVar+import Control.Concurrent import Control.Monad import Control.Monad.Except import Control.Monad.Identity import Data.ByteString (ByteString, packCStringLen, useAsCString)-import qualified Data.ByteString.Lazy.Char8 as BLC+import Data.ByteString.Lazy.Char8 qualified as BLC import Data.ByteString.Unsafe import Data.Function import Data.Text (Text)-import qualified Data.Text as T-import qualified Data.Text.Encoding as T-import qualified Data.Text.Read as T+import Data.Text qualified as T+import Data.Text.Encoding qualified as T+import Data.Text.Read qualified as T import Data.Void+import Data.Word import Foreign.C.String import Foreign.C.Types+import Foreign.ForeignPtr import Foreign.Marshal.Alloc import Foreign.Marshal.Array import Foreign.Ptr@@ -46,6 +50,7 @@ data IceSession = IceSession { isStrans :: PjIceStrans+ , _isConfig :: IceConfig , isChan :: MVar (Either [ByteString] (Flow Void ByteString)) } @@ -116,14 +121,38 @@ {#enum pj_ice_sess_role as IceSessionRole {underscoreToCase} deriving (Show, Eq) #} +data PjIceStransCfg+newtype IceConfig = IceConfig (ForeignPtr PjIceStransCfg)++foreign import ccall unsafe "pjproject.h &ice_cfg_free"+ ice_cfg_free :: FunPtr (Ptr PjIceStransCfg -> IO ())+foreign import ccall unsafe "pjproject.h ice_cfg_create"+ ice_cfg_create :: CString -> Word16 -> CString -> Word16 -> IO (Ptr PjIceStransCfg)++iceCreateConfig :: Maybe ( Text, Word16 ) -> Maybe ( Text, Word16 ) -> IO (Maybe IceConfig)+iceCreateConfig stun turn =+ maybe ($ nullPtr) (withText . fst) stun $ \cstun ->+ maybe ($ nullPtr) (withText . fst) turn $ \cturn -> do+ cfg <- ice_cfg_create cstun (maybe 0 snd stun) cturn (maybe 0 snd turn)+ if cfg == nullPtr+ then return Nothing+ else Just . IceConfig <$> newForeignPtr ice_cfg_free cfg+ {#pointer *pj_ice_strans as ^ #} -iceCreate :: IceSessionRole -> (IceSession -> IO ()) -> IO IceSession-iceCreate role cb = do+iceCreateSession :: IceConfig -> IceSessionRole -> (IceSession -> IO ()) -> IO IceSession+iceCreateSession icfg@(IceConfig fcfg) role cb = do rec sptr <- newStablePtr sess- cbptr <- newStablePtr $ cb sess+ cbptr <- newStablePtr $ do+ -- The callback may be called directly from pj_ice_strans_create or later+ -- from a different thread; make sure we use a different thread here + -- to avoid deadlock on accessing 'sess'.+ forkIO $ cb sess sess <- IceSession- <$> {#call ice_create #} (fromIntegral $ fromEnum role) (castStablePtrToPtr sptr) (castStablePtrToPtr cbptr)+ <$> (withForeignPtr fcfg $ \cfg ->+ {#call ice_create #} (castPtr cfg) (fromIntegral $ fromEnum role) (castStablePtrToPtr sptr) (castStablePtrToPtr cbptr)+ )+ <*> pure icfg <*> (newMVar $ Left []) return $ sess
src/Erebos/ICE/pjproject.c view
@@ -12,7 +12,6 @@ { pj_caching_pool cp; pj_pool_t * pool;- pj_ice_strans_cfg cfg; pj_sockaddr def_addr; } ice; @@ -31,9 +30,9 @@ fprintf(stderr, "ICE: %s: %s\n", msg, err); } -static int ice_worker_thread(void * unused)+static int ice_worker_thread(void * vcfg) {- PJ_UNUSED_ARG(unused);+ pj_ice_strans_cfg * cfg = (pj_ice_strans_cfg *) vcfg; while (true) { pj_time_val max_timeout = { 0, 0 };@@ -41,7 +40,7 @@ max_timeout.msec = 500; - pj_timer_heap_poll(ice.cfg.stun_cfg.timer_heap, &timeout);+ pj_timer_heap_poll(cfg->stun_cfg.timer_heap, &timeout); pj_assert(timeout.sec >= 0 && timeout.msec >= 0); if (timeout.msec >= 1000)@@ -50,7 +49,7 @@ if (PJ_TIME_VAL_GT(timeout, max_timeout)) timeout = max_timeout; - int c = pj_ioqueue_poll(ice.cfg.stun_cfg.ioqueue, &timeout);+ int c = pj_ioqueue_poll(cfg->stun_cfg.ioqueue, &timeout); if (c < 0) pj_thread_sleep(PJ_TIME_VAL_MSEC(timeout)); }@@ -105,7 +104,7 @@ if (done) { pthread_mutex_unlock(&mutex);- goto exit;+ return; } pj_log_set_level(1);@@ -125,49 +124,89 @@ pj_caching_pool_init(&ice.cp, NULL, 0); - pj_ice_strans_cfg_default(&ice.cfg);- ice.cfg.stun_cfg.pf = &ice.cp.factory;- ice.pool = pj_pool_create(&ice.cp.factory, "ice", 512, 512, NULL); - if (pj_timer_heap_create(ice.pool, 100,- &ice.cfg.stun_cfg.timer_heap) != PJ_SUCCESS) {- fprintf(stderr, "pj_timer_heap_create failed\n");- goto exit;+exit:+ done = true;+ pthread_mutex_unlock(&mutex);+}++pj_ice_strans_cfg * ice_cfg_create( const char * stun_server, uint16_t stun_port,+ const char * turn_server, uint16_t turn_port )+{+ ice_init();++ pj_ice_strans_cfg * cfg = malloc( sizeof(pj_ice_strans_cfg) );+ pj_ice_strans_cfg_default( cfg );++ cfg->stun_cfg.pf = &ice.cp.factory;+ if( pj_timer_heap_create( ice.pool, 100,+ &cfg->stun_cfg.timer_heap ) != PJ_SUCCESS ){+ fprintf( stderr, "pj_timer_heap_create failed\n" );+ goto fail; } - if (pj_ioqueue_create(ice.pool, 16, &ice.cfg.stun_cfg.ioqueue) != PJ_SUCCESS) {- fprintf(stderr, "pj_ioqueue_create failed\n");- goto exit;+ if( pj_ioqueue_create( ice.pool, 16, &cfg->stun_cfg.ioqueue ) != PJ_SUCCESS ){+ fprintf( stderr, "pj_ioqueue_create failed\n" );+ goto fail; } pj_thread_t * thread;- if (pj_thread_create(ice.pool, "ice", &ice_worker_thread,- NULL, 0, 0, &thread) != PJ_SUCCESS) {- fprintf(stderr, "pj_thread_create failed\n");- goto exit;+ if( pj_thread_create( ice.pool, NULL, &ice_worker_thread,+ cfg, 0, 0, &thread ) != PJ_SUCCESS ){+ fprintf( stderr, "pj_thread_create failed\n" );+ goto fail; } - ice.cfg.af = pj_AF_INET();- ice.cfg.opt.aggressive = PJ_TRUE;+ cfg->af = pj_AF_INET();+ cfg->opt.aggressive = PJ_TRUE; - ice.cfg.stun.server.ptr = "discovery1.erebosprotocol.net";- ice.cfg.stun.server.slen = strlen(ice.cfg.stun.server.ptr);- ice.cfg.stun.port = 29670;+ if( stun_server ){+ cfg->stun.server.ptr = malloc( strlen( stun_server ));+ pj_strcpy2( &cfg->stun.server, stun_server );+ if( stun_port )+ cfg->stun.port = stun_port;+ } - ice.cfg.turn.server = ice.cfg.stun.server;- ice.cfg.turn.port = ice.cfg.stun.port;- ice.cfg.turn.auth_cred.type = PJ_STUN_AUTH_CRED_STATIC;- ice.cfg.turn.auth_cred.data.static_cred.data_type = PJ_STUN_PASSWD_PLAIN;- ice.cfg.turn.conn_type = PJ_TURN_TP_UDP;+ if( turn_server ){+ cfg->turn.server.ptr = malloc( strlen( turn_server ));+ pj_strcpy2( &cfg->turn.server, turn_server );+ if( turn_port )+ cfg->turn.port = turn_port;+ cfg->turn.auth_cred.type = PJ_STUN_AUTH_CRED_STATIC;+ cfg->turn.auth_cred.data.static_cred.data_type = PJ_STUN_PASSWD_PLAIN;+ cfg->turn.conn_type = PJ_TURN_TP_UDP;+ } -exit:- done = true;- pthread_mutex_unlock(&mutex);+ return cfg;+fail:+ ice_cfg_free( cfg );+ return NULL; } -pj_ice_strans * ice_create(pj_ice_sess_role role, HsStablePtr sptr, HsStablePtr cb)+void ice_cfg_free( pj_ice_strans_cfg * cfg ) {+ if( ! cfg )+ return;++ if( cfg->turn.server.ptr )+ free( cfg->turn.server.ptr );++ if( cfg->stun.server.ptr )+ free( cfg->stun.server.ptr );++ if( cfg->stun_cfg.ioqueue )+ pj_ioqueue_destroy( cfg->stun_cfg.ioqueue );++ if( cfg->stun_cfg.timer_heap )+ pj_timer_heap_destroy( cfg->stun_cfg.timer_heap );++ free( cfg );+}++pj_ice_strans * ice_create( const pj_ice_strans_cfg * cfg, pj_ice_sess_role role,+ HsStablePtr sptr, HsStablePtr cb )+{ ice_init(); pj_ice_strans * res;@@ -182,8 +221,8 @@ .on_ice_complete = cb_on_ice_complete, }; - pj_status_t status = pj_ice_strans_create(NULL, &ice.cfg, 1,- udata, &icecb, &res);+ pj_status_t status = pj_ice_strans_create( NULL, cfg, 1,+ udata, &icecb, &res ); if (status != PJ_SUCCESS) ice_perror("error creating ice", status);@@ -358,7 +397,7 @@ return; } - pj_status_t status = pj_ice_strans_sendto(strans, 1, data, len,+ pj_status_t status = pj_ice_strans_sendto2(strans, 1, data, len, &ice.def_addr, pj_sockaddr_get_len(&ice.def_addr)); if (status != PJ_SUCCESS && status != PJ_EPENDING) ice_perror("error sending data", status);
src/Erebos/ICE/pjproject.h view
@@ -3,7 +3,12 @@ #include <pjnath.h> #include <HsFFI.h> -pj_ice_strans * ice_create(pj_ice_sess_role role, HsStablePtr sptr, HsStablePtr cb);+pj_ice_strans_cfg * ice_cfg_create( const char * stun_server, uint16_t stun_port,+ const char * turn_server, uint16_t turn_port );+void ice_cfg_free( pj_ice_strans_cfg * cfg );++pj_ice_strans * ice_create( const pj_ice_strans_cfg *, pj_ice_sess_role role,+ HsStablePtr sptr, HsStablePtr cb ); void ice_destroy(pj_ice_strans * strans); ssize_t ice_encode_session(pj_ice_strans *, char * ufrag, char * pass,
src/Erebos/Network.hs view
@@ -6,6 +6,7 @@ stopServer, getCurrentPeerList, getNextPeerChange,+ getServerAddresses, ServerOptions(..), serverIdentity, defaultServerOptions, Peer, peerServer, peerStorage,@@ -46,17 +47,17 @@ import Data.Typeable import Data.Word +import Foreign.C.Types+import Foreign.Marshal.Alloc+import Foreign.Marshal.Array import Foreign.Ptr-import Foreign.Storable+import Foreign.Storable as F import GHC.Conc.Sync (unsafeIOToSTM) 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@@ -83,6 +84,7 @@ data Server = Server { serverStorage :: Storage+ , serverOptions :: ServerOptions , serverOrigHead :: Head LocalState , serverIdentity_ :: MVar UnifiedIdentity , serverThreads :: MVar [ThreadId]@@ -229,7 +231,7 @@ return (t:ts) startServer :: ServerOptions -> Head LocalState -> (String -> IO ()) -> [SomeService] -> IO Server-startServer opt serverOrigHead logd' serverServices = do+startServer serverOptions serverOrigHead logd' serverServices = do let serverStorage = headStorage serverOrigHead serverIdentity_ <- newMVar $ headLocalIdentity serverOrigHead serverThreads <- newMVar []@@ -265,7 +267,7 @@ return sock loop sock = do- when (serverLocalDiscovery opt) $ forkServerThread server $ do+ when (serverLocalDiscovery serverOptions) $ forkServerThread server $ do announceAddreses <- fmap concat $ sequence $ [ map (SockAddrInet6 discoveryPort 0 discoveryMulticastGroup) <$> joinMulticast sock , getBroadcastAddresses discoveryPort@@ -377,7 +379,7 @@ , addrFamily = AF_INET6 , addrSocketType = Datagram }- addr:_ <- getAddrInfo (Just hints) Nothing (Just $ show $ serverPort opt)+ addr:_ <- getAddrInfo (Just hints) Nothing (Just $ show $ serverPort serverOptions) bracket (open addr) close loop forkServerThread server $ forever $ do@@ -954,17 +956,56 @@ foreign import ccall unsafe "Network/ifaddrs.h join_multicast" cJoinMulticast :: CInt -> Ptr CSize -> IO (Ptr Word32)+foreign import ccall unsafe "Network/ifaddrs.h local_addresses" cLocalAddresses :: Ptr CSize -> IO (Ptr InetAddress) foreign import ccall unsafe "Network/ifaddrs.h broadcast_addresses" cBroadcastAddresses :: IO (Ptr Word32)-foreign import ccall unsafe "stdlib.h free" cFree :: Ptr Word32 -> IO ()+foreign import ccall unsafe "stdlib.h free" cFree :: Ptr a -> IO () +data InetAddress = InetAddress { fromInetAddress :: IP.IP }++instance F.Storable InetAddress where+ sizeOf _ = sizeOf (undefined :: CInt) + 16+ alignment _ = 8++ peek ptr = (unpackFamily <$> peekByteOff ptr 0) >>= \case+ AF_INET -> InetAddress . IP.IPv4 . IP.fromHostAddress <$> peekByteOff ptr (sizeOf (undefined :: CInt))+ AF_INET6 -> InetAddress . IP.IPv6 . IP.toIPv6b . map fromIntegral <$> peekArray 16 (ptr `plusPtr` sizeOf (undefined :: CInt) :: Ptr Word8)+ _ -> fail "InetAddress: unknown family"++ poke ptr (InetAddress addr) = case addr of+ IP.IPv4 ip -> do+ pokeByteOff ptr 0 (packFamily AF_INET)+ pokeByteOff ptr (sizeOf (undefined :: CInt)) (IP.toHostAddress ip)+ IP.IPv6 ip -> do+ pokeByteOff ptr 0 (packFamily AF_INET6)+ pokeArray (ptr `plusPtr` sizeOf (undefined :: CInt) :: Ptr Word8) (map fromIntegral $ IP.fromIPv6b ip)+ 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+ if ptr == nullPtr+ then do+ return []+ else do+ count <- fromIntegral <$> peek pcount+ res <- forM [ 0 .. count - 1 ] $ \i ->+ peekElemOff ptr i+ cFree ptr+ return res++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 (IP.toSockAddr . (, serverPort serverOptions ) . fromInetAddress) res getBroadcastAddresses :: PortNumber -> IO [SockAddr] getBroadcastAddresses port = do
src/Erebos/Network/ifaddrs.c view
@@ -9,6 +9,7 @@ #ifndef _WIN32 #include <arpa/inet.h> #include <net/if.h>+#include <netinet/in.h> #include <ifaddrs.h> #include <endian.h> #include <sys/types.h>@@ -85,8 +86,73 @@ return interfaces; } +static bool copy_local_address( struct InetAddress * dst, const struct sockaddr * src )+{+ int family = src->sa_family;++ if( family == AF_INET ){+ struct in_addr * addr = & (( struct sockaddr_in * ) src)->sin_addr;+ if (! ((ntohl( addr->s_addr ) & 0xff000000) == 0x7f000000) && // loopback+ ! ((ntohl( addr->s_addr ) & 0xffff0000) == 0xa9fe0000) // link-local+ ){+ dst->family = family;+ memcpy( & dst->addr, addr, sizeof( * addr ));+ return true;+ }+ }++ if( family == AF_INET6 ){+ struct in6_addr * addr = & (( struct sockaddr_in6 * ) src)->sin6_addr;+ if (! IN6_IS_ADDR_LOOPBACK( addr ) &&+ ! IN6_IS_ADDR_LINKLOCAL( addr )+ ){+ dst->family = family;+ memcpy( & dst->addr, addr, sizeof( * addr ));+ return true;+ }+ }++ return false;+}+ #ifndef _WIN32 +struct InetAddress * local_addresses( size_t * count )+{+ struct ifaddrs * addrs;+ if( getifaddrs( &addrs ) < 0 )+ return 0;++ * count = 0;+ size_t capacity = 16;+ struct InetAddress * ret = malloc( sizeof(* ret) * capacity );++ for( struct ifaddrs * ifa = addrs; ifa; ifa = ifa->ifa_next ){+ if ( ifa->ifa_addr ){+ int family = ifa->ifa_addr->sa_family;+ if( family == AF_INET || family == AF_INET6 ){+ if( (* count) >= capacity ){+ capacity *= 2;+ struct InetAddress * nret = realloc( ret, sizeof(* ret) * capacity );+ if (nret) {+ ret = nret;+ } else {+ free( ret );+ freeifaddrs( addrs );+ return 0;+ }+ }++ if( copy_local_address( & ret[ * count ], ifa->ifa_addr ))+ (* count)++;+ }+ }+ }++ freeifaddrs(addrs);+ return ret;+}+ uint32_t * broadcast_addresses(void) { struct ifaddrs * addrs;@@ -106,6 +172,7 @@ ret = nret; } else { free(ret);+ freeifaddrs(addrs); return 0; } }@@ -124,8 +191,51 @@ #include <winsock2.h> #include <ws2tcpip.h>+#include <iptypes.h>+#include <iphlpapi.h> #pragma comment(lib, "ws2_32.lib")++struct InetAddress * local_addresses( size_t * count )+{+ * count = 0;+ struct InetAddress * ret = NULL;++ ULONG bufsize = 15000;+ IP_ADAPTER_ADDRESSES * buf = NULL;++ DWORD rv = 0;++ do {+ buf = realloc( buf, bufsize );+ rv = GetAdaptersAddresses( AF_UNSPEC, 0, NULL, buf, & bufsize );++ if( rv == ERROR_BUFFER_OVERFLOW )+ continue;+ } while (0);++ if( rv == NO_ERROR ){+ size_t capacity = 16;+ ret = malloc( sizeof( * ret ) * capacity );++ for( IP_ADAPTER_ADDRESSES * cur = (IP_ADAPTER_ADDRESSES *) buf;+ cur && (* count) < capacity;+ cur = cur->Next ){++ for( IP_ADAPTER_UNICAST_ADDRESS * curAddr = cur->FirstUnicastAddress;+ curAddr && (* count) < capacity;+ curAddr = curAddr->Next ){++ if( copy_local_address( & ret[ * count ], curAddr->Address.lpSockaddr ))+ (* count)++;+ }+ }+ }++cleanup:+ free( buf );+ return ret;+} uint32_t * broadcast_addresses(void) {
src/Erebos/Service.hs view
@@ -37,7 +37,13 @@ import Erebos.State import Erebos.Storage -class (Typeable s, Storable s, Typeable (ServiceState s), Typeable (ServiceGlobalState s)) => Service s where+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 ()
src/Erebos/Storage.hs view
@@ -567,7 +567,7 @@ True -> return ilist False -> do void $ watchDir manager (headTypePath spath tid) (const True) $ \case- Added { eventPath = fpath } | Just ihid <- HeadID <$> U.fromString (takeFileName fpath) -> do+ ev@Added {} | Just ihid <- HeadID <$> U.fromString (takeFileName (eventPath ev)) -> do loadHeadRaw st tid ihid >>= \case Just ref -> do (_, _, iwl) <- readMVar mvar@@ -841,91 +841,97 @@ loadEmpty name = maybe (throwError $ "Missing record item '"++name++"'") return =<< loadMbEmpty name loadMbEmpty :: String -> LoadRec (Maybe ())-loadMbEmpty name = (lookup (BC.pack name) <$> loadRecItems) >>= \case- Nothing -> return Nothing- Just (RecEmpty) -> return (Just ())- Just _ -> throwError $ "Expecting type int of record item '"++name++"'"+loadMbEmpty name = listToMaybe . mapMaybe p <$> loadRecItems+ where+ bname = BC.pack name+ p ( name', RecEmpty ) | name' == bname+ = Just ()+ p _ = Nothing loadInt :: Num a => String -> LoadRec a loadInt name = maybe (throwError $ "Missing record item '"++name++"'") return =<< loadMbInt name loadMbInt :: Num a => String -> LoadRec (Maybe a)-loadMbInt name = (lookup (BC.pack name) <$> loadRecItems) >>= \case- Nothing -> return Nothing- Just (RecInt x) -> return (Just $ fromInteger x)- Just _ -> throwError $ "Expecting type int of record item '"++name++"'"+loadMbInt name = listToMaybe . mapMaybe p <$> loadRecItems+ where+ bname = BC.pack name+ p ( name', RecInt x ) | name' == bname+ = Just (fromInteger x)+ p _ = Nothing loadNum :: (Real a, Fractional a) => String -> LoadRec a loadNum name = maybe (throwError $ "Missing record item '"++name++"'") return =<< loadMbNum name loadMbNum :: (Real a, Fractional a) => String -> LoadRec (Maybe a)-loadMbNum name = (lookup (BC.pack name) <$> loadRecItems) >>= \case- Nothing -> return Nothing- Just (RecNum x) -> return (Just $ fromRational x)- Just _ -> throwError $ "Expecting type number of record item '"++name++"'"+loadMbNum name = listToMaybe . mapMaybe p <$> loadRecItems+ where+ bname = BC.pack name+ p ( name', RecNum x ) | name' == bname+ = Just (fromRational x)+ p _ = Nothing loadText :: StorableText a => String -> LoadRec a loadText name = maybe (throwError $ "Missing record item '"++name++"'") return =<< loadMbText name loadMbText :: StorableText a => String -> LoadRec (Maybe a)-loadMbText name = (lookup (BC.pack name) <$> loadRecItems) >>= \case- Nothing -> return Nothing- Just (RecText x) -> Just <$> fromText x- Just _ -> throwError $ "Expecting type text of record item '"++name++"'"+loadMbText name = listToMaybe <$> loadTexts name loadTexts :: StorableText a => String -> LoadRec [a]-loadTexts name = do- items <- map snd . filter ((BC.pack name ==) . fst) <$> loadRecItems- forM items $ \case RecText x -> fromText x- _ -> throwError $ "Expecting type text of record item '"++name++"'"+loadTexts name = sequence . mapMaybe p =<< loadRecItems+ where+ bname = BC.pack name+ p ( name', RecText x ) | name' == bname+ = Just (fromText x)+ p _ = Nothing loadBinary :: BA.ByteArray a => String -> LoadRec a loadBinary name = maybe (throwError $ "Missing record item '"++name++"'") return =<< loadMbBinary name loadMbBinary :: BA.ByteArray a => String -> LoadRec (Maybe a)-loadMbBinary name = (lookup (BC.pack name) <$> loadRecItems) >>= \case- Nothing -> return Nothing- Just (RecBinary x) -> return $ Just $ BA.convert x- Just _ -> throwError $ "Expecting type binary of record item '"++name++"'"+loadMbBinary name = listToMaybe <$> loadBinaries name loadBinaries :: BA.ByteArray a => String -> LoadRec [a]-loadBinaries name = do- items <- map snd . filter ((BC.pack name ==) . fst) <$> loadRecItems- forM items $ \case RecBinary x -> return $ BA.convert x- _ -> throwError $ "Expecting type binary of record item '"++name++"'"+loadBinaries name = mapMaybe p <$> loadRecItems+ where+ bname = BC.pack name+ p ( name', RecBinary x ) | name' == bname+ = Just (BA.convert x)+ p _ = Nothing loadDate :: StorableDate a => String -> LoadRec a loadDate name = maybe (throwError $ "Missing record item '"++name++"'") return =<< loadMbDate name loadMbDate :: StorableDate a => String -> LoadRec (Maybe a)-loadMbDate name = (lookup (BC.pack name) <$> loadRecItems) >>= \case- Nothing -> return Nothing- Just (RecDate x) -> return $ Just $ fromDate x- Just _ -> throwError $ "Expecting type date of record item '"++name++"'"+loadMbDate name = listToMaybe . mapMaybe p <$> loadRecItems+ where+ bname = BC.pack name+ p ( name', RecDate x ) | name' == bname+ = Just (fromDate x)+ p _ = Nothing loadUUID :: StorableUUID a => String -> LoadRec a loadUUID name = maybe (throwError $ "Missing record iteem '"++name++"'") return =<< loadMbUUID name loadMbUUID :: StorableUUID a => String -> LoadRec (Maybe a)-loadMbUUID name = (lookup (BC.pack name) <$> loadRecItems) >>= \case- Nothing -> return Nothing- Just (RecUUID x) -> return $ Just $ fromUUID x- Just _ -> throwError $ "Expecting type UUID of record item '"++name++"'"+loadMbUUID name = listToMaybe . mapMaybe p <$> loadRecItems+ where+ bname = BC.pack name+ p ( name', RecUUID x ) | name' == bname+ = Just (fromUUID x)+ p _ = Nothing loadRawRef :: String -> LoadRec Ref loadRawRef name = maybe (throwError $ "Missing record item '"++name++"'") return =<< loadMbRawRef name loadMbRawRef :: String -> LoadRec (Maybe Ref)-loadMbRawRef name = (lookup (BC.pack name) <$> loadRecItems) >>= \case- Nothing -> return Nothing- Just (RecRef x) -> return (Just x)- Just _ -> throwError $ "Expecting type ref of record item '"++name++"'"+loadMbRawRef name = listToMaybe <$> loadRawRefs name loadRawRefs :: String -> LoadRec [Ref]-loadRawRefs name = do- items <- map snd . filter ((BC.pack name ==) . fst) <$> loadRecItems- forM items $ \case RecRef x -> return x- _ -> throwError $ "Expecting type ref of record item '"++name++"'"+loadRawRefs name = mapMaybe p <$> loadRecItems+ where+ bname = BC.pack name+ p ( name', RecRef x ) | name' == bname = Just x+ p _ = Nothing loadRef :: Storable a => String -> LoadRec a loadRef name = load <$> loadRawRef name