packages feed

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 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