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 +0/−0
- LICENSE +373/−0
- Setup.hs +6/−0
- flink-statefulfun.cabal +74/−0
- proto/Example.proto +11/−0
- proto/Kafka.proto +29/−0
- proto/RequestReply.proto +145/−0
- src/Network/Flink/ProtoServant.hs +24/−0
- src/Network/Flink/Stateful.hs +315/−0
+ 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"
+ flink-statefulfun.cabal view
@@ -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