packages feed

plexus-protocol 0.3.0.0 → 0.3.2.0

raw patch · 5 files changed

+189/−49 lines, 5 filesdep +case-insensitivedep +resource-poolPVP: major bump suggested

API removals or changes: PVP suggests a major version bump

Dependencies added: case-insensitive, resource-pool

API changes (from Hackage documentation)

- Plexus: [itemTimeoutMs] :: PlexusStreamItem -> Int
- Plexus.Types: CustomRequest :: t -> Request t
- Plexus.Types: CustomResponse :: t -> Response t
- Plexus.Types: [customRequestData] :: Request t -> t
- Plexus.Types: [customResponseData] :: Response t -> t
- Plexus.Types: [itemTimeoutMs] :: PlexusStreamItem -> Int
+ Plexus: [itemTimeout] :: PlexusStreamItem -> Maybe Int
+ Plexus: [substrateHeaders] :: SubstrateConfig -> Headers
+ Plexus.Client: [substrateHeaders] :: SubstrateConfig -> Headers
+ Plexus.Client: cookieHeader :: Text -> Headers
+ Plexus.Client: instance GHC.Classes.Ord Plexus.Client.SubstrateConfig
+ Plexus.Client.Pool: PlexusPoolConfig :: Int -> Int -> NominalDiffTime -> PlexusPoolConfig
+ Plexus.Client.Pool: [poolIdleTime] :: PlexusPoolConfig -> NominalDiffTime
+ Plexus.Client.Pool: [poolMaxResources] :: PlexusPoolConfig -> Int
+ Plexus.Client.Pool: [poolStripes] :: PlexusPoolConfig -> Int
+ Plexus.Client.Pool: createConnectionPool :: SubstrateConfig -> PlexusPoolConfig -> IO ConnectionPool
+ Plexus.Client.Pool: data PlexusPoolConfig
+ Plexus.Client.Pool: defaultPlexusPoolConfig :: PlexusPoolConfig
+ Plexus.Client.Pool: destroyConnectionPool :: ConnectionPool -> IO ()
+ Plexus.Client.Pool: instance GHC.Classes.Eq Plexus.Client.Pool.PlexusPoolConfig
+ Plexus.Client.Pool: instance GHC.Show.Show Plexus.Client.Pool.PlexusPoolConfig
+ Plexus.Client.Pool: type ConnectionPool = Pool SubstrateConnection
+ Plexus.Client.Pool: withPooledConnection :: ConnectionPool -> (SubstrateConnection -> IO a) -> IO a
+ Plexus.Types: [itemTimeout] :: PlexusStreamItem -> Maybe Int
+ Plexus.Types: type WellKnownRequest = Request Value
+ Plexus.Types: type WellKnownResponse = Response Value
- Plexus: StreamRequest :: Text -> Provenance -> Text -> Request Value -> Int -> PlexusStreamItem
+ Plexus: StreamRequest :: Text -> Provenance -> Text -> StandardRequest -> Maybe Int -> PlexusStreamItem
- Plexus: SubstrateConfig :: String -> Int -> String -> Text -> SubstrateConfig
+ Plexus: SubstrateConfig :: String -> Int -> String -> Text -> Headers -> SubstrateConfig
- Plexus: [itemRequestData] :: PlexusStreamItem -> Request Value
+ Plexus: [itemRequestData] :: PlexusStreamItem -> StandardRequest
- Plexus.Client: SubstrateConfig :: String -> Int -> String -> Text -> SubstrateConfig
+ Plexus.Client: SubstrateConfig :: String -> Int -> String -> Text -> Headers -> SubstrateConfig
- Plexus.Types: StreamRequest :: Text -> Provenance -> Text -> Request Value -> Int -> PlexusStreamItem
+ Plexus.Types: StreamRequest :: Text -> Provenance -> Text -> StandardRequest -> Maybe Int -> PlexusStreamItem
- Plexus.Types: [itemRequestData] :: PlexusStreamItem -> Request Value
+ Plexus.Types: [itemRequestData] :: PlexusStreamItem -> StandardRequest

Files

plexus-protocol.cabal view
@@ -1,6 +1,6 @@ cabal-version:      3.0 name:               plexus-protocol-version:            0.3.0.0+version:            0.3.2.0 synopsis:           Plexus RPC protocol types and client license:            MIT author:@@ -10,6 +10,7 @@ library   exposed-modules:     Plexus.Client+    Plexus.Client.Pool     Plexus.Transport     Plexus     Plexus.Types@@ -22,6 +23,7 @@     aeson-casing >= 0.2,     text >= 2.0,     bytestring >= 0.11,+    case-insensitive >= 1.2,     containers >= 0.6,     websockets >= 0.13,     network >= 3.1,@@ -30,7 +32,8 @@     streaming >= 0.2,     directory >= 1.3,     filepath >= 1.4,-    time >= 1.9+    time >= 1.9,+    resource-pool >= 0.4   hs-source-dirs:   src   default-language: GHC2021   default-extensions:
src/Plexus/Client.hs view
@@ -19,15 +19,18 @@     -- * Configuration   , SubstrateConfig(..)   , defaultConfig+  , cookieHeader   ) where -import Control.Concurrent (forkIO, threadDelay)+import Control.Concurrent (forkIO) import Control.Concurrent.Async (Async, async, cancel, wait) import Control.Concurrent.STM import Control.Exception (SomeException, catch) import Control.Monad (forever, void)+import System.Timeout (timeout) import Data.Aeson import qualified Data.Aeson.KeyMap as KM+import qualified Data.ByteString.Lazy.Char8 as LBS import Data.IORef import Data.Map.Strict (Map) import qualified Data.Map.Strict as Map@@ -35,6 +38,8 @@ import qualified Data.Text as T import Network.WebSockets (Connection) import qualified Network.WebSockets as WS+import qualified Data.ByteString.Char8 as BS8+import qualified Data.CaseInsensitive as CI import Streaming import qualified Streaming.Prelude as Str @@ -46,18 +51,24 @@   , substratePort    :: Int   , substratePath    :: String   , substrateBackend :: Text    -- ^ Backend name (e.g., "plexus")+  , substrateHeaders :: WS.Headers  -- ^ Extra HTTP upgrade headers (e.g., Cookie)   }-  deriving stock (Show, Eq)+  deriving stock (Show, Eq, Ord)  -- | Default configuration for local development (requires backend) defaultConfig :: Text -> SubstrateConfig defaultConfig backend = SubstrateConfig-  { substrateHost = "127.0.0.1"-  , substratePort = 4444-  , substratePath = "/"+  { substrateHost    = "127.0.0.1"+  , substratePort    = 4444+  , substratePath    = "/"   , substrateBackend = backend+  , substrateHeaders = []   } +-- | Build a Cookie header from a JWT token+cookieHeader :: Text -> WS.Headers+cookieHeader tok = [(CI.mk "Cookie", "access_token=" <> BS8.pack (T.unpack tok))]+ -- | Pending request with its queue, waiting for subscription ID data PendingRequest = PendingRequest   { prQueue    :: TQueue Value        -- ^ Queue to receive notifications@@ -86,7 +97,8 @@   -- Run WebSocket client in a background thread   -- Catch exceptions inside the thread to prevent them from being printed   void $ forkIO $-    (WS.runClient substrateHost substratePort substratePath $ \conn -> do+    (WS.runClientWith substrateHost substratePort substratePath+        WS.defaultConnectionOptions substrateHeaders $ \conn -> do       reader <- async $ readerLoop conn subs pendingReqs       -- Signal success       atomically $ putTMVar resultVar (Right (conn, reader))@@ -96,12 +108,11 @@       -- Signal failure (don't print, just capture)       atomically $ putTMVar resultVar (Left $ show e) -  -- Wait for connection result (with timeout)-  threadDelay 200000  -- 200ms-  mResult <- atomically $ tryTakeTMVar resultVar+  -- Wait for connection result (with 5s timeout instead of hardcoded 200ms sleep)+  mResult <- timeout 5000000 $ atomically $ takeTMVar resultVar    case mResult of-    Nothing -> error "Connection timeout"+    Nothing -> error "Connection timeout (5s)"     Just (Left err) -> error $ "Connection failed: " <> err     Just (Right (conn, reader)) ->       pure SubstrateConnection
+ src/Plexus/Client/Pool.hs view
@@ -0,0 +1,82 @@+{-# LANGUAGE RecordWildCards #-}++-- | WebSocket connection pooling for Plexus RPC+--+-- Maintains a pool of persistent WebSocket connections to avoid+-- the ~200ms overhead of creating a new connection for each RPC call.+--+-- Usage:+--+-- @+-- pool <- createConnectionPool config+-- withPooledConnection pool $ \conn -> do+--   items <- S.toList_ $ substrateRpc conn method params+--   pure items+-- @+module Plexus.Client.Pool+  ( ConnectionPool+  , createConnectionPool+  , destroyConnectionPool+  , withPooledConnection+  , defaultPlexusPoolConfig+  , PlexusPoolConfig(..)+  ) where++import Control.Exception (SomeException, catch, finally)+import Control.Monad (void)+import Data.Pool (Pool, newPool, defaultPoolConfig, setNumStripes, withResource, destroyAllResources)+import qualified Data.Pool as Pool+import Data.Time.Clock (NominalDiffTime)++import Plexus.Client+  ( SubstrateConfig+  , SubstrateConnection+  , connect+  , disconnect+  )++-- | Connection pool configuration+data PlexusPoolConfig = PlexusPoolConfig+  { poolStripes   :: Int              -- ^ Number of stripes (sub-pools)+  , poolMaxResources :: Int           -- ^ Max connections per stripe+  , poolIdleTime  :: NominalDiffTime  -- ^ How long to keep idle connections+  }+  deriving stock (Show, Eq)++-- | Default pool configuration+-- - 1 stripe (simple pool)+-- - 10 max connections+-- - 30 second idle timeout+defaultPlexusPoolConfig :: PlexusPoolConfig+defaultPlexusPoolConfig = PlexusPoolConfig+  { poolStripes = 1+  , poolMaxResources = 10+  , poolIdleTime = 30+  }++-- | Connection pool for WebSocket connections+type ConnectionPool = Pool SubstrateConnection++-- | Create a connection pool+--+-- The pool will maintain persistent WebSocket connections and reuse them+-- across multiple RPC calls, avoiding the 200ms connection setup overhead.+createConnectionPool :: SubstrateConfig -> PlexusPoolConfig -> IO ConnectionPool+createConnectionPool cfg PlexusPoolConfig{..} =+  newPool $ setNumStripes (Just poolStripes)+          $ Pool.defaultPoolConfig+              (connect cfg)                    -- Create connection+              disconnect                       -- Destroy connection+              (realToFrac poolIdleTime)        -- Idle timeout (convert to Double)+              poolMaxResources                 -- Max resources per stripe++-- | Destroy all connections in the pool+destroyConnectionPool :: ConnectionPool -> IO ()+destroyConnectionPool = destroyAllResources++-- | Use a connection from the pool+--+-- The connection is automatically returned to the pool when the action completes,+-- even if an exception is thrown.+withPooledConnection :: ConnectionPool -> (SubstrateConnection -> IO a) -> IO a+withPooledConnection = withResource
src/Plexus/Transport.hs view
@@ -27,16 +27,36 @@  import Control.Exception (SomeException, IOException, catch, fromException) import qualified Control.Exception as E+import Control.Concurrent.MVar (MVar, newMVar, modifyMVar) import Data.Aeson hiding (Error)+import Data.Map.Strict (Map)+import qualified Data.Map.Strict as Map import Data.Text (Text) import qualified Data.Text as T import qualified Streaming.Prelude as S import qualified Network.Socket as NS+import System.IO.Unsafe (unsafePerformIO)  import Plexus.Client (SubstrateConfig(..), connect, disconnect, substrateRpc, defaultConfig)+import Plexus.Client.Pool (ConnectionPool, createConnectionPool, withPooledConnection, defaultPlexusPoolConfig) import Plexus.Types (PlexusStreamItem(..), TransportError(..), Response(..), StandardResponse) import Plexus.Schema.Recursive (PluginSchema, MethodSchema, SchemaResult(..), parsePluginSchema, parseSchemaResult) +-- | Global connection pool cache+-- One pool per unique SubstrateConfig to enable connection reuse across calls+{-# NOINLINE poolCache #-}+poolCache :: MVar (Map SubstrateConfig ConnectionPool)+poolCache = unsafePerformIO (newMVar Map.empty)++-- | Get or create a connection pool for the given config+getOrCreatePool :: SubstrateConfig -> IO ConnectionPool+getOrCreatePool cfg = modifyMVar poolCache $ \pools ->+  case Map.lookup cfg pools of+    Just pool -> pure (pools, pool)+    Nothing -> do+      pool <- createConnectionPool cfg defaultPlexusPoolConfig+      pure (Map.insert cfg pool pools, pool)+ -- | Low-level RPC call with default localhost config rpcCall :: Text -> Text -> Value -> IO (Either TransportError [PlexusStreamItem]) rpcCall backend = rpcCallWith (defaultConfig backend)@@ -77,10 +97,9 @@  doCallInner :: SubstrateConfig -> Text -> Value -> IO [PlexusStreamItem] doCallInner cfg method params = do-  conn <- connect cfg-  items <- S.toList_ $ substrateRpc conn method params-  disconnect conn-  pure items+  pool <- getOrCreatePool cfg+  withPooledConnection pool $ \conn ->+    S.toList_ $ substrateRpc conn method params  -- | Streaming RPC call - invokes callback for each item as it arrives rpcCallStreaming :: SubstrateConfig -> Text -> Value -> (PlexusStreamItem -> IO ()) -> IO (Either TransportError ())@@ -91,9 +110,9 @@  doCallStreaming :: SubstrateConfig -> Text -> Value -> (PlexusStreamItem -> IO ()) -> IO () doCallStreaming cfg method params onItem = do-  conn <- connect cfg-  S.mapM_ onItem $ substrateRpc conn method params-  disconnect conn+  pool <- getOrCreatePool cfg+  withPooledConnection pool $ \conn ->+    S.mapM_ onItem $ substrateRpc conn method params  -- | Streaming method invocation invokeMethodStreaming :: SubstrateConfig -> [Text] -> Text -> Value -> (PlexusStreamItem -> IO ()) -> IO (Either TransportError ())@@ -181,7 +200,7 @@   let backend = substrateBackend cfg   let respondParams = object         [ "request_id" .= requestId-        , "response" .= response+        , "response_data" .= response         ]   result <- rpcCallWith cfg (backend <> ".respond") respondParams   case result of
src/Plexus/Types.hs view
@@ -20,9 +20,11 @@     -- * Bidirectional Types   , SelectOption(..)   , Request(..)-  , StandardRequest   , Response(..)-  , StandardResponse+  , WellKnownRequest+  , WellKnownResponse+  , StandardRequest  -- Deprecated, kept for backwards compatibility+  , StandardResponse -- Deprecated, kept for backwards compatibility      -- * Helpers   , mkSubscribeRequest@@ -337,14 +339,28 @@     <> catMaybes [("description" .=) <$> optionDescription]  -- | Parameterized request type for bidirectional communication+--+-- This type represents the well-known request types in the Plexus bidirectional protocol.+-- Unlike Rust which uses traits, Haskell uses a parameterized ADT which provides+-- similar flexibility with better type safety.+--+-- The type parameter @t@ is used for values in Prompt defaults and Select options.+-- For most cases, use 'WellKnownRequest' which sets @t = Value@. data Request t   = Confirm  { confirmMessage :: Text, confirmDefault :: Maybe Bool }   | Prompt   { promptMessage :: Text, promptDefault :: Maybe t, promptPlaceholder :: Maybe Text }   | Select   { selectMessage :: Text, selectOptions :: [SelectOption t], selectMulti :: Bool }-  | CustomRequest { customRequestData :: t }   deriving (Show, Eq, Generic) +-- | Well-known request types with JSON values+--+-- This is the recommended type for most bidirectional requests.+-- It matches the Rust @WellKnownRequest@ and TypeScript union types.+type WellKnownRequest = Request Value+ -- | Type alias for backwards compatibility: requests with JSON values+--+-- Deprecated: Use 'WellKnownRequest' instead for consistency with Rust/TypeScript. type StandardRequest = Request Value  instance (FromJSON t) => FromJSON (Request t) where@@ -362,7 +378,6 @@         <$> o .: "message"         <*> o .: "options"         <*> o .:? "multi_select" .!= False-      "custom" -> CustomRequest <$> o .: "data"       _ -> fail $ "Unknown request type: " <> T.unpack typ  instance (ToJSON t) => ToJSON (Request t) where@@ -375,21 +390,32 @@   toJSON (Select msg opts multi) = object     [ "type" .= ("select" :: Text), "message" .= msg     , "options" .= opts, "multi_select" .= multi ]-  toJSON (CustomRequest d) = object ["type" .= ("custom" :: Text), "data" .= d]  -- | Parameterized response type for bidirectional communication --+-- This type represents the well-known response types in the Plexus bidirectional protocol.+-- -- Note: The 'Value' constructor corresponds to JSON @"type":"text"@ for -- wire compatibility with the Rust implementation.+--+-- The type parameter @t@ is used for values in Prompt responses and Select responses.+-- For most cases, use 'WellKnownResponse' which sets @t = Value@. data Response t   = Confirmed { confirmedValue :: Bool }   | Value     { responseValue :: t }    -- ^ wire tag: "text"   | Selected  { selectedValues :: [t] }-  | CustomResponse { customResponseData :: t }   | Cancelled   deriving (Show, Eq, Generic) +-- | Well-known response types with JSON values+--+-- This is the recommended type for most bidirectional responses.+-- It matches the Rust @WellKnownResponse@ and TypeScript union types.+type WellKnownResponse = Response Value+ -- | Type alias for backwards compatibility: responses with JSON values+--+-- Deprecated: Use 'WellKnownResponse' instead for consistency with Rust/TypeScript. type StandardResponse = Response Value  instance (FromJSON t) => FromJSON (Response t) where@@ -399,16 +425,14 @@       "confirmed" -> Confirmed <$> o .: "value"       "text"      -> Value <$> o .: "value"       "selected"  -> Selected <$> o .: "values"-      "custom"    -> CustomResponse <$> o .: "data"       "cancelled" -> pure Cancelled       _ -> fail $ "Unknown response type: " <> T.unpack typ  instance (ToJSON t) => ToJSON (Response t) where-  toJSON (Confirmed v)       = object ["type" .= ("confirmed" :: Text), "value" .= v]-  toJSON (Value v)           = object ["type" .= ("text" :: Text), "value" .= v]-  toJSON (Selected vs)       = object ["type" .= ("selected" :: Text), "values" .= vs]-  toJSON (CustomResponse d)  = object ["type" .= ("custom" :: Text), "data" .= d]-  toJSON Cancelled           = object ["type" .= ("cancelled" :: Text)]+  toJSON (Confirmed v)  = object ["type" .= ("confirmed" :: Text), "value" .= v]+  toJSON (Value v)      = object ["type" .= ("text" :: Text), "value" .= v]+  toJSON (Selected vs)  = object ["type" .= ("selected" :: Text), "values" .= vs]+  toJSON Cancelled      = object ["type" .= ("cancelled" :: Text)]  -- ============================================================================ -- Plexus Stream Items@@ -445,17 +469,17 @@       , itemError       :: Text       , itemRecoverable :: Bool       }-  | StreamDone-      { itemPlexusHash :: Text-      , itemProvenance :: Provenance-      }   | StreamRequest       { itemPlexusHash  :: Text       , itemProvenance  :: Provenance       , itemRequestId   :: Text-      , itemRequestData :: Request Value-      , itemTimeoutMs   :: Int+      , itemRequestData :: StandardRequest+      , itemTimeout     :: Maybe Int       }+  | StreamDone+      { itemPlexusHash :: Text+      , itemProvenance :: Provenance+      }   deriving stock (Show, Eq, Generic)  instance FromJSON PlexusStreamItem where@@ -488,11 +512,11 @@           "error" -> StreamError hash prov             <$> o .: "message"             <*> o .:? "recoverable" .!= False-          "done" -> pure $ StreamDone hash prov           "request" -> StreamRequest hash prov             <$> o .: "request_id"             <*> o .: "request"-            <*> o .: "timeout_ms"+            <*> o .:? "timeout"+          "done" -> pure $ StreamDone hash prov           _ -> fail $ "Unknown event type: " <> T.unpack typ        -- Legacy format: flat structure@@ -518,13 +542,13 @@             <$> o .: "provenance"             <*> o .: "error"             <*> o .: "recoverable"-          "done" -> StreamDone hash-            <$> o .: "provenance"           "request" -> StreamRequest hash             <$> o .: "provenance"             <*> o .: "request_id"             <*> o .: "request"-            <*> o .: "timeout_ms"+            <*> o .:? "timeout"+          "done" -> StreamDone hash+            <$> o .: "provenance"           _ -> fail $ "Unknown event type: " <> T.unpack typ  instance ToJSON PlexusStreamItem where@@ -555,16 +579,17 @@     , "message" .= err     , "recoverable" .= rec     ]-  toJSON (StreamDone hash prov) = object-    [ "type" .= ("done" :: Text)-    , "metadata" .= StreamMetadata prov hash 0-    ]-  toJSON (StreamRequest hash prov reqId reqData timeout) = object+  toJSON (StreamRequest hash prov reqId reqData timeout) = object $     [ "type" .= ("request" :: Text)     , "metadata" .= StreamMetadata prov hash 0     , "request_id" .= reqId     , "request" .= reqData-    , "timeout_ms" .= timeout+    ] <> catMaybes+    [ ("timeout" .=) <$> timeout+    ]+  toJSON (StreamDone hash prov) = object+    [ "type" .= ("done" :: Text)+    , "metadata" .= StreamMetadata prov hash 0     ]  -- | Create a subscription request