diff --git a/CHANGELOG.md b/CHANGELOG.md
new file mode 100644
--- /dev/null
+++ b/CHANGELOG.md
diff --git a/LICENSE b/LICENSE
new file mode 100644
--- /dev/null
+++ b/LICENSE
@@ -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.
diff --git a/Setup.hs b/Setup.hs
new file mode 100644
--- /dev/null
+++ b/Setup.hs
@@ -0,0 +1,6 @@
+-- import           Distribution.Simple
+-- main = defaultMain
+
+import Data.ProtoLens.Setup
+
+main = defaultMainGeneratingProtos "proto"
diff --git a/flink-statefulfun.cabal b/flink-statefulfun.cabal
new file mode 100644
--- /dev/null
+++ b/flink-statefulfun.cabal
@@ -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
diff --git a/proto/Example.proto b/proto/Example.proto
new file mode 100644
--- /dev/null
+++ b/proto/Example.proto
@@ -0,0 +1,11 @@
+syntax = "proto3";
+
+package example;
+
+message GreeterRequest {
+  string name = 1;
+}
+
+message GreeterResponse {
+  string greeting = 1;
+}
diff --git a/proto/Kafka.proto b/proto/Kafka.proto
new file mode 100644
--- /dev/null
+++ b/proto/Kafka.proto
@@ -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;
+}
diff --git a/proto/RequestReply.proto b/proto/RequestReply.proto
new file mode 100644
--- /dev/null
+++ b/proto/RequestReply.proto
@@ -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;
+    }
+}
+
+
diff --git a/src/Network/Flink/ProtoServant.hs b/src/Network/Flink/ProtoServant.hs
new file mode 100644
--- /dev/null
+++ b/src/Network/Flink/ProtoServant.hs
@@ -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
diff --git a/src/Network/Flink/Stateful.hs b/src/Network/Flink/Stateful.hs
new file mode 100644
--- /dev/null
+++ b/src/Network/Flink/Stateful.hs
@@ -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
