|
| 1 | +module Comm |
| 2 | + ( KadOp(..), sendLookup, newUid, parseHeader, |
| 3 | + ) where |
| 4 | + |
| 5 | +import Network.Socket |
| 6 | +import Network.BSD |
| 7 | + |
| 8 | +import Data.Word |
| 9 | +import Data.List (genericDrop,unfoldr) |
| 10 | +import Data.Char(chr,ord) |
| 11 | +import Control.Monad(forM, ap, liftM) |
| 12 | +import Control.Monad.State |
| 13 | + |
| 14 | +import KTable |
| 15 | +import Globals |
| 16 | + |
| 17 | +-- Header: version - 1 byte |
| 18 | +-- optype - 1 byte |
| 19 | +-- uid - 8 bytes |
| 20 | +-- |
| 21 | +-- Node Lookup: node id - 20 bytes |
| 22 | + |
| 23 | +data PeerHandle = PeerHandle { pSocket :: Socket, pAddress :: SockAddr } |
| 24 | + |
| 25 | +-- Client functions to send operations to other nodes. |
| 26 | +-- |
| 27 | + |
| 28 | +sendLookup peers nid lookupId = forM peers (\p -> do |
| 29 | + msgId <- liftIO newUid |
| 30 | + msg <- buildLookupMsg nid msgId |
| 31 | + newWaitingReply p NodeLookupOp msgId lookupId |
| 32 | + liftIO $ sendToPeer msg p ) |
| 33 | + |
| 34 | + where buildLookupMsg nid msgId = liftM (++ toCharArray nid 160) $ buildHeader NodeLookupOp msgId |
| 35 | + |
| 36 | +sendToPeer msg peer = do |
| 37 | + phandle <- openPeerHandle (host peer) (port peer) |
| 38 | + sendstr phandle msg |
| 39 | + closePeerHandle phandle |
| 40 | + |
| 41 | + where sendstr _ [] = return () |
| 42 | + sendstr phandle omsg = do |
| 43 | + sent <- sendTo (pSocket phandle) omsg (pAddress phandle) |
| 44 | + sendstr phandle (genericDrop sent omsg) |
| 45 | + |
| 46 | +-- Server functions to handle calls coming from other nodes |
| 47 | +-- |
| 48 | + |
| 49 | +serverDispatch addr msg = do |
| 50 | + let (hdr, rest) = parseHeader msg |
| 51 | + -- ignoring what we can't handle |
| 52 | + if msgVersion hdr /= 1 |
| 53 | + then return () |
| 54 | + else do wait <- waitingReply (msgUid hdr) (msgOp hdr) (toPeer addr) |
| 55 | + case wait of |
| 56 | + Nothing -> return () |
| 57 | + opId -> dispatchOp (msgOp hdr) opId rest |
| 58 | + |
| 59 | + where dispatchOp NodeLookupOp opId msg = do |
| 60 | + rl <- runningLookup opId |
| 61 | + case rl of |
| 62 | + Nothing -> return () |
| 63 | + Just rl -> return () |
| 64 | + dispatchOp _ _ _ = return () |
| 65 | + |
| 66 | +data Header = Header { msgVersion ::Int, msgOp ::KadOp, msgUid ::Word64 } |
| 67 | + deriving Show |
| 68 | + |
| 69 | +-- TODO handle things like empty messages, adding error handling |
| 70 | +parseHeader msg = runState parseHeader' msg |
| 71 | + where parseHeader' = do |
| 72 | + ver <- parseVersion |
| 73 | + op <- parseOpType |
| 74 | + uid <- parseUid |
| 75 | + return $ Header ver op uid |
| 76 | + |
| 77 | + parseVersion = consuming 1 (ord . head) |
| 78 | + parseOpType = consuming 1 (\x -> |
| 79 | + case head x of |
| 80 | + '\SOH' -> PingOp |
| 81 | + '\STX' -> NodeLookupOp |
| 82 | + _ -> UnknownOp ) |
| 83 | + parseUid = consuming 8 fromCharArray |
| 84 | + |
| 85 | + consuming n fn = do |
| 86 | + str <- get |
| 87 | + let (v, r) = splitAt n str |
| 88 | + put r |
| 89 | + return $ fn v |
| 90 | + |
| 91 | +localServer port handlerFn = withSocketsDo $ do |
| 92 | + addrinfos <- getAddrInfo (Just (defaultHints {addrFlags = [AI_PASSIVE]})) Nothing (Just port) |
| 93 | + let serveraddr = head addrinfos |
| 94 | + sock <- socket (addrFamily serveraddr) Datagram defaultProtocol |
| 95 | + bindSocket sock (addrAddress serveraddr) |
| 96 | + procMessages sock |
| 97 | + |
| 98 | + -- Loops forever processing incoming data |
| 99 | + where procMessages sock = do |
| 100 | + (msg, _, addr) <- recvFrom sock 1024 |
| 101 | + handlerFn addr msg |
| 102 | + procMessages sock |
| 103 | + |
| 104 | +-- Utility functions to serialize messages |
| 105 | +-- |
| 106 | + |
| 107 | +buildHeader optype msgId = do |
| 108 | + return $ buildHeader' msgId optype |
| 109 | + |
| 110 | + where buildHeader':: (Integral a) => a -> KadOp -> String |
| 111 | + buildHeader' uid optype = (chr 1) : optypeStr : toCharArray uid 64 |
| 112 | + optypeStr = case optype of |
| 113 | + PingOp -> chr 1 |
| 114 | + NodeLookupOp -> chr 2 |
| 115 | + |
| 116 | +toCharArray:: (Integral a) => a -> Int -> [Char] |
| 117 | +toCharArray num depth = map chr (toBytes num depth) |
| 118 | + |
| 119 | +-- Converts a string to its numeric value by considering that each character is a |
| 120 | +-- byte in a n byte number |
| 121 | +fromCharArray :: (Num t) => [Char] -> t |
| 122 | +fromCharArray str = fst $ foldl (\(acc, exp) ch -> |
| 123 | + (acc + fromIntegral (ord ch) * 2^exp, exp+1)) (0,0) (reverse str) |
| 124 | + |
| 125 | +-- Decomposes a number into its successive byte values or, said differently, converts |
| 126 | +-- a number to base 2^8. |
| 127 | +toBytes :: (Integral t, Integral a) => a -> t -> [Int] |
| 128 | +toBytes num depth = unfoldr byteMod (num,depth) |
| 129 | + where byteMod (num,exp) = let (d,m) = divMod num (2^exp) |
| 130 | + in if exp < 0 then Nothing else Just (fromIntegral d ::Int, (m,exp-8)) |
| 131 | + |
| 132 | +toPeer addr = |
| 133 | + case addr of |
| 134 | + SockAddrInet port host -> Peer (show host) (show port) (-1) |
| 135 | + SockAddrInet6 port _ host _ -> Peer (show host) (show port) (-1) |
| 136 | + SockAddrUnix str -> error $ "Dont know how to translate " ++ str |
| 137 | + |
| 138 | +openPeerHandle hostname port = do |
| 139 | + addrinfos <- getAddrInfo Nothing (Just hostname) (Just port) |
| 140 | + let serveraddr = head addrinfos |
| 141 | + sock <- socket (addrFamily serveraddr) Datagram defaultProtocol |
| 142 | + return $ PeerHandle sock (addrAddress serveraddr) |
| 143 | + |
| 144 | +closePeerHandle phandle = sClose (pSocket phandle) |
| 145 | + |
| 146 | +-- main = do |
| 147 | +-- params <- getArgs |
| 148 | +-- if head params == "client" |
| 149 | +-- then do |
| 150 | +-- ph <- openPeerHandle "localhost" "10000" |
| 151 | +-- sendPing ph |
| 152 | +-- closePeerHandle ph |
| 153 | +-- else localServer "10000" (\addr msg -> putStrLn $ msg ++ " / " ++ (show addr)) |
0 commit comments