packages feed

flink-statefulfun (empty) → 0.1.0.0

raw patch · 9 files changed

+977/−0 lines, 9 filesdep +basedep +bytestringdep +containersbuild-type:Customsetup-changed

Dependencies added: base, bytestring, containers, either, http-media, http-types, lens-family, mtl, proto-lens, proto-lens-protobuf-types, proto-lens-runtime, servant, servant-server, text, wai, warp

Files

+ CHANGELOG.md view
+ LICENSE view
@@ -0,0 +1,373 @@+Mozilla Public License Version 2.0+==================================++1. Definitions+--------------++1.1. "Contributor"+    means each individual or legal entity that creates, contributes to+    the creation of, or owns Covered Software.++1.2. "Contributor Version"+    means the combination of the Contributions of others (if any) used+    by a Contributor and that particular Contributor's Contribution.++1.3. "Contribution"+    means Covered Software of a particular Contributor.++1.4. "Covered Software"+    means Source Code Form to which the initial Contributor has attached+    the notice in Exhibit A, the Executable Form of such Source Code+    Form, and Modifications of such Source Code Form, in each case+    including portions thereof.++1.5. "Incompatible With Secondary Licenses"+    means++    (a) that the initial Contributor has attached the notice described+        in Exhibit B to the Covered Software; or++    (b) that the Covered Software was made available under the terms of+        version 1.1 or earlier of the License, but not also under the+        terms of a Secondary License.++1.6. "Executable Form"+    means any form of the work other than Source Code Form.++1.7. "Larger Work"+    means a work that combines Covered Software with other material, in+    a separate file or files, that is not Covered Software.++1.8. "License"+    means this document.++1.9. "Licensable"+    means having the right to grant, to the maximum extent possible,+    whether at the time of the initial grant or subsequently, any and+    all of the rights conveyed by this License.++1.10. "Modifications"+    means any of the following:++    (a) any file in Source Code Form that results from an addition to,+        deletion from, or modification of the contents of Covered+        Software; or++    (b) any new file in Source Code Form that contains any Covered+        Software.++1.11. "Patent Claims" of a Contributor+    means any patent claim(s), including without limitation, method,+    process, and apparatus claims, in any patent Licensable by such+    Contributor that would be infringed, but for the grant of the+    License, by the making, using, selling, offering for sale, having+    made, import, or transfer of either its Contributions or its+    Contributor Version.++1.12. "Secondary License"+    means either the GNU General Public License, Version 2.0, the GNU+    Lesser General Public License, Version 2.1, the GNU Affero General+    Public License, Version 3.0, or any later versions of those+    licenses.++1.13. "Source Code Form"+    means the form of the work preferred for making modifications.++1.14. "You" (or "Your")+    means an individual or a legal entity exercising rights under this+    License. For legal entities, "You" includes any entity that+    controls, is controlled by, or is under common control with You. For+    purposes of this definition, "control" means (a) the power, direct+    or indirect, to cause the direction or management of such entity,+    whether by contract or otherwise, or (b) ownership of more than+    fifty percent (50%) of the outstanding shares or beneficial+    ownership of such entity.++2. License Grants and Conditions+--------------------------------++2.1. Grants++Each Contributor hereby grants You a world-wide, royalty-free,+non-exclusive license:++(a) under intellectual property rights (other than patent or trademark)+    Licensable by such Contributor to use, reproduce, make available,+    modify, display, perform, distribute, and otherwise exploit its+    Contributions, either on an unmodified basis, with Modifications, or+    as part of a Larger Work; and++(b) under Patent Claims of such Contributor to make, use, sell, offer+    for sale, have made, import, and otherwise transfer either its+    Contributions or its Contributor Version.++2.2. Effective Date++The licenses granted in Section 2.1 with respect to any Contribution+become effective for each Contribution on the date the Contributor first+distributes such Contribution.++2.3. Limitations on Grant Scope++The licenses granted in this Section 2 are the only rights granted under+this License. No additional rights or licenses will be implied from the+distribution or licensing of Covered Software under this License.+Notwithstanding Section 2.1(b) above, no patent license is granted by a+Contributor:++(a) for any code that a Contributor has removed from Covered Software;+    or++(b) for infringements caused by: (i) Your and any other third party's+    modifications of Covered Software, or (ii) the combination of its+    Contributions with other software (except as part of its Contributor+    Version); or++(c) under Patent Claims infringed by Covered Software in the absence of+    its Contributions.++This License does not grant any rights in the trademarks, service marks,+or logos of any Contributor (except as may be necessary to comply with+the notice requirements in Section 3.4).++2.4. Subsequent Licenses++No Contributor makes additional grants as a result of Your choice to+distribute the Covered Software under a subsequent version of this+License (see Section 10.2) or under the terms of a Secondary License (if+permitted under the terms of Section 3.3).++2.5. Representation++Each Contributor represents that the Contributor believes its+Contributions are its original creation(s) or it has sufficient rights+to grant the rights to its Contributions conveyed by this License.++2.6. Fair Use++This License is not intended to limit any rights You have under+applicable copyright doctrines of fair use, fair dealing, or other+equivalents.++2.7. Conditions++Sections 3.1, 3.2, 3.3, and 3.4 are conditions of the licenses granted+in Section 2.1.++3. Responsibilities+-------------------++3.1. Distribution of Source Form++All distribution of Covered Software in Source Code Form, including any+Modifications that You create or to which You contribute, must be under+the terms of this License. You must inform recipients that the Source+Code Form of the Covered Software is governed by the terms of this+License, and how they can obtain a copy of this License. You may not+attempt to alter or restrict the recipients' rights in the Source Code+Form.++3.2. Distribution of Executable Form++If You distribute Covered Software in Executable Form then:++(a) such Covered Software must also be made available in Source Code+    Form, as described in Section 3.1, and You must inform recipients of+    the Executable Form how they can obtain a copy of such Source Code+    Form by reasonable means in a timely manner, at a charge no more+    than the cost of distribution to the recipient; and++(b) You may distribute such Executable Form under the terms of this+    License, or sublicense it under different terms, provided that the+    license for the Executable Form does not attempt to limit or alter+    the recipients' rights in the Source Code Form under this License.++3.3. Distribution of a Larger Work++You may create and distribute a Larger Work under terms of Your choice,+provided that You also comply with the requirements of this License for+the Covered Software. If the Larger Work is a combination of Covered+Software with a work governed by one or more Secondary Licenses, and the+Covered Software is not Incompatible With Secondary Licenses, this+License permits You to additionally distribute such Covered Software+under the terms of such Secondary License(s), so that the recipient of+the Larger Work may, at their option, further distribute the Covered+Software under the terms of either this License or such Secondary+License(s).++3.4. Notices++You may not remove or alter the substance of any license notices+(including copyright notices, patent notices, disclaimers of warranty,+or limitations of liability) contained within the Source Code Form of+the Covered Software, except that You may alter any license notices to+the extent required to remedy known factual inaccuracies.++3.5. Application of Additional Terms++You may choose to offer, and to charge a fee for, warranty, support,+indemnity or liability obligations to one or more recipients of Covered+Software. However, You may do so only on Your own behalf, and not on+behalf of any Contributor. You must make it absolutely clear that any+such warranty, support, indemnity, or liability obligation is offered by+You alone, and You hereby agree to indemnify every Contributor for any+liability incurred by such Contributor as a result of warranty, support,+indemnity or liability terms You offer. You may include additional+disclaimers of warranty and limitations of liability specific to any+jurisdiction.++4. Inability to Comply Due to Statute or Regulation+---------------------------------------------------++If it is impossible for You to comply with any of the terms of this+License with respect to some or all of the Covered Software due to+statute, judicial order, or regulation then You must: (a) comply with+the terms of this License to the maximum extent possible; and (b)+describe the limitations and the code they affect. Such description must+be placed in a text file included with all distributions of the Covered+Software under this License. Except to the extent prohibited by statute+or regulation, such description must be sufficiently detailed for a+recipient of ordinary skill to be able to understand it.++5. Termination+--------------++5.1. The rights granted under this License will terminate automatically+if You fail to comply with any of its terms. However, if You become+compliant, then the rights granted under this License from a particular+Contributor are reinstated (a) provisionally, unless and until such+Contributor explicitly and finally terminates Your grants, and (b) on an+ongoing basis, if such Contributor fails to notify You of the+non-compliance by some reasonable means prior to 60 days after You have+come back into compliance. Moreover, Your grants from a particular+Contributor are reinstated on an ongoing basis if such Contributor+notifies You of the non-compliance by some reasonable means, this is the+first time You have received notice of non-compliance with this License+from such Contributor, and You become compliant prior to 30 days after+Your receipt of the notice.++5.2. If You initiate litigation against any entity by asserting a patent+infringement claim (excluding declaratory judgment actions,+counter-claims, and cross-claims) alleging that a Contributor Version+directly or indirectly infringes any patent, then the rights granted to+You by any and all Contributors for the Covered Software under Section+2.1 of this License shall terminate.++5.3. In the event of termination under Sections 5.1 or 5.2 above, all+end user license agreements (excluding distributors and resellers) which+have been validly granted by You or Your distributors under this License+prior to termination shall survive termination.++************************************************************************+*                                                                      *+*  6. Disclaimer of Warranty                                           *+*  -------------------------                                           *+*                                                                      *+*  Covered Software is provided under this License on an "as is"       *+*  basis, without warranty of any kind, either expressed, implied, or  *+*  statutory, including, without limitation, warranties that the       *+*  Covered Software is free of defects, merchantable, fit for a        *+*  particular purpose or non-infringing. The entire risk as to the     *+*  quality and performance of the Covered Software is with You.        *+*  Should any Covered Software prove defective in any respect, You     *+*  (not any Contributor) assume the cost of any necessary servicing,   *+*  repair, or correction. This disclaimer of warranty constitutes an   *+*  essential part of this License. No use of any Covered Software is   *+*  authorized under this License except under this disclaimer.         *+*                                                                      *+************************************************************************++************************************************************************+*                                                                      *+*  7. Limitation of Liability                                          *+*  --------------------------                                          *+*                                                                      *+*  Under no circumstances and under no legal theory, whether tort      *+*  (including negligence), contract, or otherwise, shall any           *+*  Contributor, or anyone who distributes Covered Software as          *+*  permitted above, be liable to You for any direct, indirect,         *+*  special, incidental, or consequential damages of any character      *+*  including, without limitation, damages for lost profits, loss of    *+*  goodwill, work stoppage, computer failure or malfunction, or any    *+*  and all other commercial damages or losses, even if such party      *+*  shall have been informed of the possibility of such damages. This   *+*  limitation of liability shall not apply to liability for death or   *+*  personal injury resulting from such party's negligence to the       *+*  extent applicable law prohibits such limitation. Some               *+*  jurisdictions do not allow the exclusion or limitation of           *+*  incidental or consequential damages, so this exclusion and          *+*  limitation may not apply to You.                                    *+*                                                                      *+************************************************************************++8. Litigation+-------------++Any litigation relating to this License may be brought only in the+courts of a jurisdiction where the defendant maintains its principal+place of business and such litigation shall be governed by laws of that+jurisdiction, without reference to its conflict-of-law provisions.+Nothing in this Section shall prevent a party's ability to bring+cross-claims or counter-claims.++9. Miscellaneous+----------------++This License represents the complete agreement concerning the subject+matter hereof. If any provision of this License is held to be+unenforceable, such provision shall be reformed only to the extent+necessary to make it enforceable. Any law or regulation which provides+that the language of a contract shall be construed against the drafter+shall not be used to construe this License against a Contributor.++10. Versions of the License+---------------------------++10.1. New Versions++Mozilla Foundation is the license steward. Except as provided in Section+10.3, no one other than the license steward has the right to modify or+publish new versions of this License. Each version will be given a+distinguishing version number.++10.2. Effect of New Versions++You may distribute the Covered Software under the terms of the version+of the License under which You originally received the Covered Software,+or under the terms of any subsequent version published by the license+steward.++10.3. Modified Versions++If you create software not governed by this License, and you want to+create a new license for such software, you may create and use a+modified version of this License if you rename the license and remove+any references to the name of the license steward (except to note that+such modified license differs from this License).++10.4. Distributing Source Code Form that is Incompatible With Secondary+Licenses++If You choose to distribute Source Code Form that is Incompatible With+Secondary Licenses under the terms of this version of the License, the+notice described in Exhibit B of this License must be attached.++Exhibit A - Source Code Form License Notice+-------------------------------------------++  This Source Code Form is subject to the terms of the Mozilla Public+  License, v. 2.0. If a copy of the MPL was not distributed with this+  file, You can obtain one at http://mozilla.org/MPL/2.0/.++If it is not possible or desirable to put the notice in a particular+file, then You may include the notice in a location (such as a LICENSE+file in a relevant directory) where a recipient would be likely to look+for such a notice.++You may add additional accurate notices of copyright ownership.++Exhibit B - "Incompatible With Secondary Licenses" Notice+---------------------------------------------------------++  This Source Code Form is "Incompatible With Secondary Licenses", as+  defined by the Mozilla Public License, v. 2.0.
+ Setup.hs view
@@ -0,0 +1,6 @@+-- import           Distribution.Simple+-- main = defaultMain++import Data.ProtoLens.Setup++main = defaultMainGeneratingProtos "proto"
@@ -0,0 +1,74 @@+cabal-version:       2.4+-- Initial package description 'policy.cabal' generated by 'cabal init'.+-- For further documentation, see http://haskell.org/cabal/users-guide/++name:                flink-statefulfun+version:             0.1.0.0+synopsis:            Flink stateful functions SDK+description:+    Typeclasses for serving Flink stateful functions+    from Haskell.+-- bug-reports:+license:             MPL-2.0+license-file:        LICENSE+author:              Timothy Bess+maintainer:          tdbgamer@gmail.com+-- copyright:+-- category:++-- build-type:          Simple+build-type: Custom++extra-source-files:  CHANGELOG.md+                     proto/**/*.proto++custom-setup+  setup-depends: base, Cabal, proto-lens-setup++library+  ghc-options:         -Wall+  build-tool-depends:  proto-lens-protoc:proto-lens-protoc+  hs-source-dirs:      src+  autogen-modules:     Proto.RequestReply Proto.RequestReply_Fields+                       Proto.Kafka Proto.Kafka_Fields+  other-modules:       Network.Flink.ProtoServant+  exposed-modules:     Network.Flink.Stateful+                       Proto.RequestReply Proto.RequestReply_Fields+                       Proto.Kafka Proto.Kafka_Fields+  -- exposed-modules:+  -- other-extensions:+  build-depends:       base >= 4.0 && < 4.15+                     , text >= 1.0 && < 1.3+                     , bytestring >= 0.10 && < 0.11+                     , either >= 5 && < 5.1+                     , containers >= 0.5 && < 0.7+                     , wai >= 3 && < 3.3+                     , warp >= 3 && < 3.4+                     , http-types >= 0.10 && < 0.13+                     , http-media >= 0.7 && < 0.9+                     , mtl >= 2 && < 2.3+                     , lens-family >= 2 && < 2.2+                     , proto-lens >= 0.5 && < 0.8+                     , proto-lens-runtime >= 0.5 && < 0.8+                     , proto-lens-protobuf-types >= 0.5 && < 0.8+                     , servant >= 0.16 && < 0.19+                     , servant-server >= 0.16 && < 0.19++  default-extensions:  OverloadedStrings+                       QuasiQuotes+                       StrictData+                       GeneralizedNewtypeDeriving+                       FlexibleContexts+                       FlexibleInstances+                       LambdaCase+                       TupleSections+                       DataKinds+                       TypeOperators+                       DeriveGeneric+                       DeriveFunctor+                       ExistentialQuantification+                       MultiParamTypeClasses+                       FunctionalDependencies+                       ConstraintKinds++  default-language:    Haskell2010
+ proto/Example.proto view
@@ -0,0 +1,11 @@+syntax = "proto3";++package example;++message GreeterRequest {+  string name = 1;+}++message GreeterResponse {+  string greeting = 1;+}
+ proto/Kafka.proto view
@@ -0,0 +1,29 @@+/*+ * Licensed to the Apache Software Foundation (ASF) under one+ * or more contributor license agreements.  See the NOTICE file+ * distributed with this work for additional information+ * regarding copyright ownership.  The ASF licenses this file+ * to you under the Apache License, Version 2.0 (the+ * "License"); you may not use this file except in compliance+ * with the License.  You may obtain a copy of the License at+ *+ *     http://www.apache.org/licenses/LICENSE-2.0+ *+ * Unless required by applicable law or agreed to in writing, software+ * distributed under the License is distributed on an "AS IS" BASIS,+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.+ * See the License for the specific language governing permissions and+ * limitations under the License.+ */++syntax = "proto3";++package org.apache.flink.statefun.flink.io;+option java_package = "org.apache.flink.statefun.flink.io.generated";+option java_multiple_files = true;++message KafkaProducerRecord {+    string key = 1;+    bytes value_bytes = 2;+    string topic = 3;+}
+ proto/RequestReply.proto view
@@ -0,0 +1,145 @@+/*+ * Licensed to the Apache Software Foundation (ASF) under one+ * or more contributor license agreements.  See the NOTICE file+ * distributed with this work for additional information+ * regarding copyright ownership.  The ASF licenses this file+ * to you under the Apache License, Version 2.0 (the+ * "License"); you may not use this file except in compliance+ * with the License.  You may obtain a copy of the License at+ *+ *     http://www.apache.org/licenses/LICENSE-2.0+ *+ * Unless required by applicable law or agreed to in writing, software+ * distributed under the License is distributed on an "AS IS" BASIS,+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.+ * See the License for the specific language governing permissions and+ * limitations under the License.+ */++syntax = "proto3";++package org.apache.flink.statefun.flink.core.polyglot;+option java_package = "org.apache.flink.statefun.flink.core.polyglot.generated";+option java_multiple_files = true;++import "google/protobuf/any.proto";++// -------------------------------------------------------------------------------------------------------------------+// Common message definitions+// -------------------------------------------------------------------------------------------------------------------++// An Address is the unique identity of an individual StatefulFunction, containing+// a function's type and an unique identifier within the type. The function's+// type denotes the "class" of function to invoke, while the unique identifier addresses the+// invocation to a specific function instance.+message Address {+    string namespace = 1;+    string type = 2;+    string id = 3;+}++// -------------------------------------------------------------------------------------------------------------------+// Messages sent to a Remote Function  +// -------------------------------------------------------------------------------------------------------------------++// The following section contains all the message types that are sent +// from Flink to a remote function.+message ToFunction {+    // PersistedValue represents a PersistedValue's value that is managed by Flink on behalf of a remote function. +    message PersistedValue {+        // The unique name of the persisted state.+        string state_name = 1;+        // The serialized state value+        bytes state_value = 2;+    }++    // Invocation represents a remote function call, it associated with an (optional) return address,+    // and an argument. +    message Invocation {+        // The address of the function that requested the invocation (possibly absent)+        Address caller = 1;+        // The invocation argument (aka the message sent to the target function)+        google.protobuf.Any argument = 2;+    }++    // InvocationBatchRequest represents a request to invoke a remote function. It is always associated with a target+    // address (the function to invoke), a list of eager state values.+    message InvocationBatchRequest {+        // The address of the function to invoke+        Address target = 1;+        // A list of PersistedValues that were registered as an eager state.+        repeated PersistedValue state = 2;+        // A non empty (at least one) list of invocations+        repeated Invocation invocations = 3;+    }++    oneof request {+        InvocationBatchRequest invocation = 100;+    }+}++// -------------------------------------------------------------------------------------------------------------------+// Messages sent from a Remote Function  +// -------------------------------------------------------------------------------------------------------------------++// The following section contains messages sent from a remote function back to Flink. +message FromFunction {+    // MutatePersistedValueCommand represents a command sent from a remote function to Flink,+    // requesting a change to a persisted value.+    message PersistedValueMutation {+        enum MutationType {+            DELETE = 0;+            MODIFY = 1;+        }+        MutationType mutation_type = 1;+        string state_name = 2;+        bytes state_value = 3;+    }++    // Invocation represents a remote function call, it associated with a (mandatory) target address,+    // and an argument. +    message Invocation {+        // The target function to invoke +        Address target = 1;+        // The invocation argument (aka the message sent to the target function)+        google.protobuf.Any argument = 2;+    }++    // DelayedInvocation represents a delayed remote function call with a target address, an argument+    // and a delay in milliseconds, after which this message to be sent.+    message DelayedInvocation {+        // the amount of milliseconds to wait before sending this message+        int64 delay_in_ms = 1;+        // the target address to send this message to+        Address target = 2;+        // the invocation argument+        google.protobuf.Any argument = 3;+    }++    // EgressMessage an argument to forward to an egress.+    // An egress is identified by a namespace and type (see EgressIdentifier SDK class).+    // The argument is a google.protobuf.Any+    message EgressMessage {+        // The target egress namespace+        string egress_namespace = 1;+        // The target egress type+        string egress_type = 2;+        // egress argument+        google.protobuf.Any argument = 3;+    }++    // InvocationResponse represents a result of an org.apache.flink.statefun.flink.core.polyglot.ToFunction.InvocationBatchRequest+    // it contains a list of state mutation to preform as a result of computing this batch, and a list of outgoing messages.+    message InvocationResponse {+        repeated PersistedValueMutation state_mutations = 1;+        repeated Invocation outgoing_messages = 2;+        repeated DelayedInvocation delayed_invocations = 3;+        repeated EgressMessage outgoing_egresses = 4;+    }++    oneof response {+        InvocationResponse invocation_result = 100;+    }+}++
+ src/Network/Flink/ProtoServant.hs view
@@ -0,0 +1,24 @@+module Network.Flink.ProtoServant where++import qualified Data.ByteString.Lazy as BS+import qualified Data.List.NonEmpty as NE+import qualified Data.ProtoLens as P+import Network.HTTP.Media ((//))+import qualified Servant.API.ContentTypes as S++data Proto++instance S.Accept Proto where+  contentTypes _ =+    NE.fromList+      [ "application" // "octet-stream",+        "application" // "protobuf",+        "application" // "x-protobuf",+        "application" // "vnd.google.protobuf"+      ]++instance P.Message m => S.MimeRender Proto m where+  mimeRender _ = BS.fromStrict . P.encodeMessage++instance P.Message m => S.MimeUnrender Proto m where+  mimeUnrender _ = P.decodeMessage . BS.toStrict
+ src/Network/Flink/Stateful.hs view
@@ -0,0 +1,315 @@+module Network.Flink.Stateful+  ( StatefulFunc+      ( insideCtx,+        getCtx,+        setCtx,+        modifyCtx,+        sendMsg,+        sendMsgDelay,+        sendEgressMsg+      ),+    makeConcrete,+    createApp,+    flinkServer,+    flinkApi,+    kafkaRecord,+    Function,+    FlinkState (..),+    FunctionTable+  )+where++import Control.Monad.Except+import Control.Monad.Reader+import Control.Monad.State (MonadState, StateT (..), gets, modify)+import Data.ByteString (ByteString)+import qualified Data.ByteString.Lazy.Char8 as BSL+import Data.Either.Combinators (mapLeft)+import Data.Foldable (Foldable (toList))+import Data.Map (Map)+import qualified Data.Map as Map+import Data.Maybe (listToMaybe)+import Data.ProtoLens (Message, defMessage, encodeMessage)+import Data.ProtoLens.Any (UnpackError)+import qualified Data.ProtoLens.Any as Any+import Data.ProtoLens.Prism+import Data.Sequence (Seq)+import qualified Data.Sequence as Seq+import Data.Text (Text)+import Data.Text.Lazy (fromStrict)+import qualified Data.Text.Lazy.Encoding as T+import Lens.Family2+import Network.Flink.ProtoServant (Proto)+import Proto.Google.Protobuf.Any (Any)+import qualified Proto.Kafka as Kafka+import qualified Proto.Kafka_Fields as Kafka+import Proto.RequestReply (FromFunction, ToFunction)+import qualified Proto.RequestReply as PR+import qualified Proto.RequestReply_Fields as PR+import Servant++--- | Table of stateful functions `(functionNamespace, functionType) -> (initialState, function)+type FunctionTable = Map (Text, Text) (ByteString, ByteString -> Env -> PR.ToFunction'InvocationBatchRequest -> IO (Either FlinkError (FunctionState ByteString)))++data Env = Env+  { envFunctionNamespace :: Text,+    envFunctionType :: Text,+    envFunctionId :: Text+  }+  deriving (Show)++data FunctionState ctx = FunctionState+  { functionStateCtx :: ctx,+    functionStateMutated :: Bool,+    functionStateInvocations :: Seq PR.FromFunction'Invocation,+    functionStateDelayedInvocations :: Seq PR.FromFunction'DelayedInvocation,+    functionStateEgressMessages :: Seq PR.FromFunction'EgressMessage+  }+  deriving (Show, Functor)++newState :: a -> FunctionState a+newState initialCtx = FunctionState initialCtx False mempty mempty mempty++-- | Monad stack used for the execution of a Flink stateful function+-- Don't reference this directly in your code if possible+newtype Function s a = Function {runFunction :: ExceptT FlinkError (StateT (FunctionState s) (ReaderT Env IO)) a}+  deriving (Monad, Applicative, Functor, MonadState (FunctionState s), MonadError FlinkError, MonadIO, MonadReader Env)++-- | Provides functions for Flink state SerDe+class FlinkState s where+  -- | decodes Flink state types from strict 'ByteString's+  decodeState :: ByteString -> Either String s+  -- | encodes Flink state types to strict 'ByteString's+  encodeState :: s -> ByteString++instance FlinkState () where+  decodeState _ = pure ()+  encodeState _ = ""++{-| Used to represent all Flink stateful function capabilities.++Contexts are received from Flink and deserialized into `s`+all modifications to state are shipped back to Flink at the end of the+batch to be persisted.++Message passing is also queued up and passed back at the end of the current+batch.++Example of a stateless function (done by setting `s` to `()`) that adds one+to a number and puts the protobuf response on Kafka via an egress message:++@+adder :: StatefulFunc () m => AdderRequest -> m ()+adder msg = sendEgressMsg ("adder", "added") (kafkaRecord "added" name added)+  where+    num = msg ^. AdderRequest.num+    added = defMessage & AdderResponse.num .~ (num + 1)+@++Example of a stateful function:++@+newtype GreeterState = GreeterState+  { greeterStateCount :: Int+  }+  deriving (Generic, Show, ToJSON, FromJSON)++instance FlinkState GreeterState where+  decodeState = eitherDecode . BSL.fromStrict+  encodeState = BSL.toStrict . Data.Aeson.encode++counter :: StatefulFunc GreeterState m => EX.GreeterRequest -> m ()+counter msg = do+  newCount \<\- (+ 1) \<$> insideCtx greeterStateCount+  let respMsg = "Saw " <> T.unpack name <> " " <> show newCount <> " time(s)"++  sendEgressMsg ("greeting", "greets") (kafkaRecord "greets" name $ response (T.pack respMsg))+  modifyCtx (\old -> old {greeterStateCount = newCount})+  where+    name = msg ^. EX.name+    response :: Text -> EX.GreeterResponse+    response greeting =+      defMessage+        & EX.greeting .~ greeting+@++This will respond to each event by counting how many times it has been called for the name it was passed.+The final state is taken and sent back to Flink. Failures of any kind will cause state to rollback to+previous values seamlessly without double counting.+-}+class MonadIO m => StatefulFunc s m | m -> s where+  -- Internal+  setInitialCtx :: s -> m ()++  -- Public+  insideCtx :: (s -> a) -> m a+  getCtx :: m s+  setCtx :: s -> m ()+  modifyCtx :: (s -> s) -> m ()+  sendMsg ::+    Message a =>+    -- | Function address (namespace, type, id)+    (Text, Text, Text) ->+    -- | protobuf message to send+    a ->+    m ()+  sendMsgDelay ::+    Message a =>+    -- | Function address (namespace, type, id)+    (Text, Text, Text) ->+    -- | delay before message send+    Int ->+    -- | protobuf message to send+    a ->+    m ()+  sendEgressMsg ::+    Message a =>+    -- | egress address (namespace, type)+    (Text, Text) ->+    -- | protobuf message to send (should be a Kafka or Kinesis protobuf record)+    a ->+    m ()++instance (FlinkState s) => StatefulFunc s (Function s) where+  setInitialCtx ctx = modify (\old -> old {functionStateCtx = ctx})++  insideCtx func = func <$> getCtx+  getCtx = gets functionStateCtx+  setCtx new = modify (\old -> old {functionStateCtx = new, functionStateMutated = True})+  modifyCtx mutator = mutator <$> getCtx >>= setCtx+  sendMsg (namespace, funcType, id') msg = do+    invocations <- gets functionStateInvocations+    modify (\old -> old {functionStateInvocations = invocations Seq.:|> invocation})+    where+      target :: PR.Address+      target =+        defMessage+          & PR.namespace .~ namespace+          & PR.type' .~ funcType+          & PR.id .~ id'+      invocation :: PR.FromFunction'Invocation+      invocation =+        defMessage+          & PR.target .~ target+          & PR.argument .~ Any.pack msg+  sendMsgDelay (namespace, funcType, id') delay msg = do+    invocations <- gets functionStateDelayedInvocations+    modify (\old -> old {functionStateDelayedInvocations = invocations Seq.:|> invocation})+    where+      target :: PR.Address+      target =+        defMessage+          & PR.namespace .~ namespace+          & PR.type' .~ funcType+          & PR.id .~ id'+      invocation :: PR.FromFunction'DelayedInvocation+      invocation =+        defMessage+          & PR.delayInMs .~ fromIntegral delay+          & PR.target .~ target+          & PR.argument .~ Any.pack msg+  sendEgressMsg (namespace, egressType) msg = do+    egresses <- gets functionStateEgressMessages+    modify (\old -> old {functionStateEgressMessages = egresses Seq.:|> egressMsg})+    where+      egressMsg :: PR.FromFunction'EgressMessage+      egressMsg =+        defMessage+          & PR.egressNamespace .~ namespace+          & PR.egressType .~ egressType+          & PR.argument .~ Any.pack msg++data FlinkError+  = MissingInvocationBatch+  | ProtoUnpackError UnpackError+  | ProtoDeserializeError String+  | StateDecodeError String+  | NoSuchFunction (Text, Text)+  deriving (Show, Eq)++invoke :: (FlinkState s, StatefulFunc s m, MonadError FlinkError m, Message a, MonadReader Env m) => (a -> m b) -> Any -> m b+invoke f input = f =<< liftEither (mapLeft ProtoUnpackError $ Any.unpack input)++-- | Takes a function taking an abstract state/message type and converts it to take concrete 'ByteString's+-- This allows each function in the 'FunctionTable' to take its own individual type of state and just expose+-- a function accepting 'ByteString' to the library code.+makeConcrete :: (FlinkState s, Message a) => (a -> Function s ()) -> ByteString -> Env -> PR.ToFunction'InvocationBatchRequest -> IO (Either FlinkError (FunctionState ByteString))+makeConcrete func initialContext env invocationBatch = runExceptT $ do+  deserializedContext <- liftEither $ mapLeft StateDecodeError $ decodeState initialContext+  (err, finalState) <- liftIO $ runner (newState deserializedContext)+  liftEither err+  return $ encodeState <$> finalState+  where+    runner state = runReaderT (runStateT (runExceptT $ runFunction runWithCtx) state) env+    runWithCtx = do+      defaultCtx <- gets functionStateCtx+      let initialCtx = getInitialCtx defaultCtx (listToMaybe $ invocationBatch ^. PR.state)+      case initialCtx of+        Left err -> throwError err+        Right ctx -> setInitialCtx ctx+      mapM_ (invoke func) ((^. PR.argument) <$> invocationBatch ^. PR.invocations)++    getInitialCtx def pv = handleEmptyState def $ (^. PR.stateValue) <$> pv+    handleEmptyState def state' = case state' of+      Just "" -> return def+      Just other -> mapLeft StateDecodeError $ decodeState other+      Nothing -> return def++createFlinkResp :: FunctionState ByteString -> FromFunction+createFlinkResp (FunctionState state mutated invocations delayedInvocations egresses) =+  defMessage & PR.invocationResult+    .~ ( defMessage+           & PR.stateMutations .~ toList stateMutations+           & PR.outgoingMessages .~ toList invocations+           & PR.delayedInvocations .~ toList delayedInvocations+           & PR.outgoingEgresses .~ toList egresses+       )+  where+    stateMutations :: [PR.FromFunction'PersistedValueMutation]+    stateMutations =+      [ defMessage+          & PR.mutationType .~ PR.FromFunction'PersistedValueMutation'MODIFY+          & PR.stateName .~ "flink_state"+          & PR.stateValue .~ state+        | mutated+      ]++type FlinkApi =+  "statefun" :> ReqBody '[Proto] ToFunction :> Post '[Proto] FromFunction++flinkApi :: Proxy FlinkApi+flinkApi = Proxy++createApp :: FunctionTable -> Application+createApp funcs = serve flinkApi (flinkServer funcs)++flinkServer :: FunctionTable -> Server FlinkApi+flinkServer functions toFunction = do+  batch <- getBatch toFunction+  ((initialCtx, function), (namespace, type', id')) <- findFunc (batch ^. PR.target)+  result <- liftIO $ function initialCtx (Env namespace type' id') batch+  finalState <- liftEither $ mapLeft flinkErrToServant result+  return $ createFlinkResp finalState+  where+    getBatch input = maybe (throwError $ flinkErrToServant MissingInvocationBatch) return (input ^? PR.maybe'request . _Just . PR._ToFunction'Invocation')+    findFunc addr = do+      res <- maybe (throwError $ flinkErrToServant $ NoSuchFunction (namespace, type')) return (Map.lookup (namespace, type') functions)+      return (res, address)+      where+        address@(namespace, type', _) = (addr ^. PR.namespace, addr ^. PR.type', addr ^. PR.id)++flinkErrToServant :: FlinkError -> ServerError+flinkErrToServant err = case err of+  MissingInvocationBatch -> err400 {errBody = "Invocation batch missing"}+  ProtoUnpackError unpackErr -> err400 {errBody = "Failed to unpack protobuf Any " <> BSL.pack (show unpackErr)}+  ProtoDeserializeError protoErr -> err400 {errBody = "Could not deserialize protobuf " <> BSL.pack protoErr}+  StateDecodeError decodeErr -> err400 {errBody = "Invalid JSON " <> BSL.pack decodeErr}+  NoSuchFunction (namespace, type') -> err400 {errBody = "No such function " <> T.encodeUtf8 (fromStrict namespace) <> T.encodeUtf8 (fromStrict type')}++-- | Takes a `topic`, `key`, and protobuf `value` to construct 'KafkaProducerRecord's for egress+kafkaRecord :: (Message v) => Text -> Text -> v -> Kafka.KafkaProducerRecord+kafkaRecord topic k v =+  defMessage+    & Kafka.topic .~ topic+    & Kafka.key .~ k+    & Kafka.valueBytes .~ encodeMessage v