packages feed

prodapi-pg (empty) → 0.1.0.0

raw patch · 6 files changed

+618/−0 lines, 6 filesdep +aesondep +basedep +binary

Dependencies added: aeson, base, binary, bytestring, clock, postgresql-simple, prodapi-core, prodapi-pg, text

Files

+ CHANGELOG.md view
@@ -0,0 +1,5 @@+# Revision history for prodapi-pg++## 0.1.0.0 -- YYYY-mm-dd++* First version. Released on an unsuspecting world.
+ LICENSE view
@@ -0,0 +1,202 @@++                                 Apache License+                           Version 2.0, January 2004+                        http://www.apache.org/licenses/++   TERMS AND CONDITIONS FOR USE, REPRODUCTION, AND DISTRIBUTION++   1. Definitions.++      "License" shall mean the terms and conditions for use, reproduction,+      and distribution as defined by Sections 1 through 9 of this document.++      "Licensor" shall mean the copyright owner or entity authorized by+      the copyright owner that is granting the License.++      "Legal Entity" shall mean the union of the acting entity and all+      other entities that control, are controlled by, or are under common+      control with that entity. For the purposes of this definition,+      "control" means (i) the power, direct or indirect, to cause the+      direction or management of such entity, whether by contract or+      otherwise, or (ii) ownership of fifty percent (50%) or more of the+      outstanding shares, or (iii) beneficial ownership of such entity.++      "You" (or "Your") shall mean an individual or Legal Entity+      exercising permissions granted by this License.++      "Source" form shall mean the preferred form for making modifications,+      including but not limited to software source code, documentation+      source, and configuration files.++      "Object" form shall mean any form resulting from mechanical+      transformation or translation of a Source form, including but+      not limited to compiled object code, generated documentation,+      and conversions to other media types.++      "Work" shall mean the work of authorship, whether in Source or+      Object form, made available under the License, as indicated by a+      copyright notice that is included in or attached to the work+      (an example is provided in the Appendix below).++      "Derivative Works" shall mean any work, whether in Source or Object+      form, that is based on (or derived from) the Work and for which the+      editorial revisions, annotations, elaborations, or other modifications+      represent, as a whole, an original work of authorship. For the purposes+      of this License, Derivative Works shall not include works that remain+      separable from, or merely link (or bind by name) to the interfaces of,+      the Work and Derivative Works thereof.++      "Contribution" shall mean any work of authorship, including+      the original version of the Work and any modifications or additions+      to that Work or Derivative Works thereof, that is intentionally+      submitted to Licensor for inclusion in the Work by the copyright owner+      or by an individual or Legal Entity authorized to submit on behalf of+      the copyright owner. For the purposes of this definition, "submitted"+      means any form of electronic, verbal, or written communication sent+      to the Licensor or its representatives, including but not limited to+      communication on electronic mailing lists, source code control systems,+      and issue tracking systems that are managed by, or on behalf of, the+      Licensor for the purpose of discussing and improving the Work, but+      excluding communication that is conspicuously marked or otherwise+      designated in writing by the copyright owner as "Not a Contribution."++      "Contributor" shall mean Licensor and any individual or Legal Entity+      on behalf of whom a Contribution has been received by Licensor and+      subsequently incorporated within the Work.++   2. Grant of Copyright License. Subject to the terms and conditions of+      this License, each Contributor hereby grants to You a perpetual,+      worldwide, non-exclusive, no-charge, royalty-free, irrevocable+      copyright license to reproduce, prepare Derivative Works of,+      publicly display, publicly perform, sublicense, and distribute the+      Work and such Derivative Works in Source or Object form.++   3. Grant of Patent License. Subject to the terms and conditions of+      this License, each Contributor hereby grants to You a perpetual,+      worldwide, non-exclusive, no-charge, royalty-free, irrevocable+      (except as stated in this section) patent license to make, have made,+      use, offer to sell, sell, import, and otherwise transfer the Work,+      where such license applies only to those patent claims licensable+      by such Contributor that are necessarily infringed by their+      Contribution(s) alone or by combination of their Contribution(s)+      with the Work to which such Contribution(s) was submitted. If You+      institute patent litigation against any entity (including a+      cross-claim or counterclaim in a lawsuit) alleging that the Work+      or a Contribution incorporated within the Work constitutes direct+      or contributory patent infringement, then any patent licenses+      granted to You under this License for that Work shall terminate+      as of the date such litigation is filed.++   4. Redistribution. You may reproduce and distribute copies of the+      Work or Derivative Works thereof in any medium, with or without+      modifications, and in Source or Object form, provided that You+      meet the following conditions:++      (a) You must give any other recipients of the Work or+          Derivative Works a copy of this License; and++      (b) You must cause any modified files to carry prominent notices+          stating that You changed the files; and++      (c) You must retain, in the Source form of any Derivative Works+          that You distribute, all copyright, patent, trademark, and+          attribution notices from the Source form of the Work,+          excluding those notices that do not pertain to any part of+          the Derivative Works; and++      (d) If the Work includes a "NOTICE" text file as part of its+          distribution, then any Derivative Works that You distribute must+          include a readable copy of the attribution notices contained+          within such NOTICE file, excluding those notices that do not+          pertain to any part of the Derivative Works, in at least one+          of the following places: within a NOTICE text file distributed+          as part of the Derivative Works; within the Source form or+          documentation, if provided along with the Derivative Works; or,+          within a display generated by the Derivative Works, if and+          wherever such third-party notices normally appear. The contents+          of the NOTICE file are for informational purposes only and+          do not modify the License. You may add Your own attribution+          notices within Derivative Works that You distribute, alongside+          or as an addendum to the NOTICE text from the Work, provided+          that such additional attribution notices cannot be construed+          as modifying the License.++      You may add Your own copyright statement to Your modifications and+      may provide additional or different license terms and conditions+      for use, reproduction, or distribution of Your modifications, or+      for any such Derivative Works as a whole, provided Your use,+      reproduction, and distribution of the Work otherwise complies with+      the conditions stated in this License.++   5. Submission of Contributions. Unless You explicitly state otherwise,+      any Contribution intentionally submitted for inclusion in the Work+      by You to the Licensor shall be under the terms and conditions of+      this License, without any additional terms or conditions.+      Notwithstanding the above, nothing herein shall supersede or modify+      the terms of any separate license agreement you may have executed+      with Licensor regarding such Contributions.++   6. Trademarks. This License does not grant permission to use the trade+      names, trademarks, service marks, or product names of the Licensor,+      except as required for reasonable and customary use in describing the+      origin of the Work and reproducing the content of the NOTICE file.++   7. Disclaimer of Warranty. Unless required by applicable law or+      agreed to in writing, Licensor provides the Work (and each+      Contributor provides its Contributions) on an "AS IS" BASIS,+      WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or+      implied, including, without limitation, any warranties or conditions+      of TITLE, NON-INFRINGEMENT, MERCHANTABILITY, or FITNESS FOR A+      PARTICULAR PURPOSE. You are solely responsible for determining the+      appropriateness of using or redistributing the Work and assume any+      risks associated with Your exercise of permissions under this License.++   8. Limitation of Liability. In no event and under no legal theory,+      whether in tort (including negligence), contract, or otherwise,+      unless required by applicable law (such as deliberate and grossly+      negligent acts) or agreed to in writing, shall any Contributor be+      liable to You for damages, including any direct, indirect, special,+      incidental, or consequential damages of any character arising as a+      result of this License or out of the use or inability to use the+      Work (including but not limited to damages for loss of goodwill,+      work stoppage, computer failure or malfunction, or any and all+      other commercial damages or losses), even if such Contributor+      has been advised of the possibility of such damages.++   9. Accepting Warranty or Additional Liability. While redistributing+      the Work or Derivative Works thereof, You may choose to offer,+      and charge a fee for, acceptance of support, warranty, indemnity,+      or other liability obligations and/or rights consistent with this+      License. However, in accepting such obligations, You may act only+      on Your own behalf and on Your sole responsibility, not on behalf+      of any other Contributor, and only if You agree to indemnify,+      defend, and hold each Contributor harmless for any liability+      incurred by, or claims asserted against, such Contributor by reason+      of your accepting any such warranty or additional liability.++   END OF TERMS AND CONDITIONS++   APPENDIX: How to apply the Apache License to your work.++      To apply the Apache License to your work, attach the following+      boilerplate notice, with the fields enclosed by brackets "[]"+      replaced with your own identifying information. (Don't include+      the brackets!)  The text should be enclosed in the appropriate+      comment syntax for the file format. We also recommend that a+      file or class name and description of purpose be included on the+      same "printed page" as the copyright notice for easier+      identification within third-party archives.++   Copyright [yyyy] [name of copyright owner]++   Licensed 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.
+ prodapi-pg.cabal view
@@ -0,0 +1,108 @@+cabal-version:      3.0+-- The cabal-version field refers to the version of the .cabal specification,+-- and can be different from the cabal-install (the tool) version and the+-- Cabal (the library) version you are using. As such, the Cabal (the library)+-- version used must be equal or greater than the version stated in this field.+-- Starting from the specification version 2.2, the cabal-version field must be+-- the first thing in the .cabal file.++-- Initial package description 'prodapi-pg' generated by+-- 'cabal init'. For further documentation, see:+--   http://haskell.org/cabal/users-guide+--+-- The name of the package.+name:               prodapi-pg++-- The package version.+-- See the Haskell package versioning policy (PVP) for standards+-- guiding when and how versions should be incremented.+-- https://pvp.haskell.org+-- PVP summary:     +-+------- breaking API changes+--                  | | +----- non-breaking API additions+--                  | | | +--- code changes with no API change+version:            0.1.0.0++-- A short (one-line) description of the package.+-- synopsis:+synopsis: ProdAPI wrappers for postgresql-simple.++-- A longer description of the package.+description: Provides ProdAPI-style tracing around postgresql-simple.++-- The license under which the package is released.+license:            Apache-2.0++-- The file containing the license text.+license-file:       LICENSE++-- The package author(s).+author:             Lucas DiCioccio++-- An email address to which users can send suggestions, bug reports, and patches.+maintainer:         lucas@dicioccio.fr++-- A copyright notice.+-- copyright:+category:           Prod+build-type:         Simple++-- Extra doc files to be distributed with the package, such as a CHANGELOG or a README.+extra-doc-files:    CHANGELOG.md++-- Extra source files to be distributed with the package, such as examples, or a tutorial module.+-- extra-source-files:++common warnings+    ghc-options: -Wall++library+    -- Import common warning flags.+    import:           warnings++    -- Modules exported by the library.+    exposed-modules:  Prod.Pg.DatabaseUtils+                      Prod.Pg.TaskQueue++    -- Modules included in this library but not exported.+    -- other-modules:++    -- LANGUAGE extensions used by modules in this package.+    -- other-extensions:++    -- Other library packages from which modules are imported.+    build-depends:    base >=4.19.1.0 && <5,+                      prodapi-core >= 0.1.0 && < 0.2,+                      aeson >= 2.2.3.0 && < 2.3,+                      binary >= 0.8.9.3 && < 0.9,+                      bytestring >= 0.12.1.0 && < 0.13,+                      clock >=0.8.4 && < 0.9,+                      postgresql-simple >= 0.7.0.0 && < 0.8,+                      text >= 2.1.2 && < 2.2 ++    -- Directories containing source files.+    hs-source-dirs:   src++    -- Base language which the package is written in.+    default-language: Haskell2010++test-suite prodapi-pg-test+    -- Import common warning flags.+    import:           warnings++    -- Base language which the package is written in.+    default-language: Haskell2010++    -- The interface type and version of the test suite.+    type:             exitcode-stdio-1.0++    -- Directories containing source files.+    hs-source-dirs:   test++    -- The entrypoint to the test suite.+    main-is:          Main.hs++    -- Test dependencies.+    build-depends:+        base ^>=4.19.1.0,+        prodapi-pg+
+ src/Prod/Pg/DatabaseUtils.hs view
@@ -0,0 +1,113 @@+{-# LANGUAGE OverloadedRecordDot #-}+{-# LANGUAGE OverloadedStrings #-}+{-# LANGUAGE QuasiQuotes #-}+{-# LANGUAGE ScopedTypeVariables #-}++module Prod.Pg.DatabaseUtils where++import qualified Data.Binary.Builder as Builder+import qualified Data.ByteString as ByteString+import qualified Data.ByteString.Char8 as C8+import Data.Int (Int64)+import Database.PostgreSQL.Simple (Connection, FromRow, Only (..), ToRow)+import qualified Database.PostgreSQL.Simple as Sql+import Database.PostgreSQL.Simple.SqlQQ (sql)+import Database.PostgreSQL.Simple.ToField (Action (..), ToField (..))+import qualified Database.PostgreSQL.Simple.Types as Sql+import Prod.Tracer (Tracer (..), contramap)+import qualified System.Clock as Clock++-------------------------------------------------------------------------------++data PGConnectionTrace+    = RunQuery ByteString.ByteString+    | DoneQuery Clock.TimeSpec Clock.TimeSpec Int ByteString.ByteString+    | RunExec Sql.Query+    | DoneExec Clock.TimeSpec Clock.TimeSpec Int64 Sql.Query++instance Show PGConnectionTrace where+    show (RunQuery bs) = "`" <> C8.unpack bs <> "`"+    show (DoneQuery t0 t1 k _) = mconcat ["n-rows=", show k, "; delay=", show (Clock.diffTimeSpec t1 t0)]+    show (RunExec q) = "`" <> show q <> "`"+    show (DoneExec t0 t1 k _) = mconcat ["result=", show k, "; delay=", show (Clock.diffTimeSpec t1 t0)]++pgExec ::+    (ToRow q) =>+    Tracer IO PGConnectionTrace ->+    Connection ->+    Sql.Query ->+    q ->+    IO Int64+pgExec tracer conn q args = do+    runTracer tracer (RunExec q)+    t0 <- Clock.getTime Clock.Monotonic+    ret <- Sql.execute conn q args+    t1 <- Clock.getTime Clock.Monotonic+    runTracer tracer (DoneExec t0 t1 ret q)+    pure ret++pgQuery ::+    (ToRow q, FromRow r) =>+    Tracer IO PGConnectionTrace ->+    Connection ->+    Sql.Query ->+    q ->+    IO [r]+pgQuery tracer conn q args = do+    formatted <- Sql.formatQuery conn q args+    runTracer tracer (RunQuery formatted)+    t0 <- Clock.getTime Clock.Monotonic+    ret <- Sql.query conn q args+    t1 <- Clock.getTime Clock.Monotonic+    runTracer tracer (DoneQuery t0 t1 (length ret) formatted)+    pure ret++-------------------------------------------------------------------------------+data Trace+    = RefreshMview MviewName ConcurrentlyOrNot PGConnectionTrace+    | RefillMtable MtableFillingFunction PGConnectionTrace+    deriving (Show)++-------------------------------------------------------------------------------+newtype MviewName = MviewName ByteString.ByteString+    deriving (Show)++instance ToField MviewName where+    toField (MviewName n) = Plain (Builder.fromByteString n)++data ConcurrentlyOrNot+    = Concurrently+    | Blocking+    deriving (Show, Eq)++refreshMview ::+    Tracer IO Trace ->+    IO Connection ->+    MviewName ->+    ConcurrentlyOrNot ->+    IO ()+refreshMview tracer mkConn name c = do+    conn <- mkConn+    _ <- pgExec (contramap (RefreshMview name c) tracer) conn q (Only name) :: IO (Int64)+    pure ()+  where+    q = case c of+        Blocking -> [sql|REFRESH MATERIALIZED VIEW ? |]+        Concurrently -> [sql|REFRESH MATERIALIZED VIEW CONCURRENTLY ? |]++-------------------------------------------------------------------------------+newtype MtableFillingFunction = MtableFillingFunction ByteString.ByteString+    deriving (Show)++execFillMTable ::+    Tracer IO Trace ->+    IO Connection ->+    MtableFillingFunction ->+    IO ()+execFillMTable tracer mkConn selectfunc@(MtableFillingFunction sqlfun) = do+    conn <- mkConn+    _ <- pgQuery (contramap (RefillMtable selectfunc) tracer) conn q () :: IO [(Only ())]+    pure ()+  where+    q :: Sql.Query+    q = Sql.Query ("select " <> sqlfun <> "()")
+ src/Prod/Pg/TaskQueue.hs view
@@ -0,0 +1,186 @@+{-# LANGUAGE OverloadedRecordDot #-}+{-# LANGUAGE OverloadedStrings #-}+{-# LANGUAGE QuasiQuotes #-}+{-# LANGUAGE ScopedTypeVariables #-}++-- todo: concurrency control with+-- \* monotonic version number and dead/alive checks+-- \* association to locks+module Prod.Pg.TaskQueue where++import qualified Data.Aeson as Aeson+import Data.Foldable (for_)+import Data.Int (Int64)+import qualified Data.Maybe as Maybe+import qualified Data.Text as Text+import Data.Typeable (Typeable)+import Database.PostgreSQL.Simple (Connection, Only (..))+import Database.PostgreSQL.Simple.Newtypes as SqlNewtypes+import Database.PostgreSQL.Simple.SqlQQ (sql)+import Prod.Stepper (Delayable (..), StepIO)+import qualified Prod.Stepper as Stepper+import Prod.Tracer (Tracer (..), contramap)++-------------------------------------------------------------------------------+import Prod.Pg.DatabaseUtils (PGConnectionTrace, pgQuery)++-------------------------------------------------------------------------------+type TaskId = Int64++data TaskHandle task+    = TaskHandle+    { taskId :: TaskId+    , taskDefinition :: task+    , markStarted :: Connection -> IO ()+    , markSuspended :: Connection -> IO ()+    , markFinished :: Connection -> IO ()+    }++data PGTrace+    = Enqueue PGConnectionTrace+    | Enqueued TaskId+    | NextTask PGConnectionTrace+    | GotNextTask (Maybe TaskId)+    | UpdateStatus Status TaskId PGConnectionTrace+    | UpdatedTask Status (Maybe TaskId)+    deriving (Show)++enqueue ::+    forall task.+    (Aeson.ToJSON task) =>+    Tracer IO PGTrace ->+    Connection ->+    task ->+    IO ()+enqueue tracer conn task = do+    let args = (Only $ SqlNewtypes.Aeson task)+    xs :: [(Only TaskId)] <- pgQuery (contramap Enqueue tracer) conn q args+    for_ xs $ \(Only x) ->+        runTracer tracer (Enqueued x)+  where+    q =+        [sql|INSERT INTO task_queue(status, priority, payload)+               VALUES ('new', 50, ?)+               RETURNING (id) |]++nextTask ::+    forall task.+    (Typeable task, Aeson.FromJSON task) =>+    Tracer IO PGTrace ->+    Connection ->+    IO (Maybe (TaskHandle task))+nextTask tracer conn = do+    let args = (Only ("6 hours" :: Text.Text))+    xs :: [(TaskId, SqlNewtypes.Aeson task)] <- pgQuery (contramap NextTask tracer) conn q args+    let xTask = Maybe.listToMaybe xs+    runTracer tracer (GotNextTask $ fmap fst xTask)+    return $ do+        (tId, payload) <- xTask+        pure $+            TaskHandle+                tId+                (SqlNewtypes.getAeson payload)+                (updateTaskStatus tracer "started" tId)+                (updateTaskStatus tracer "suspended" tId)+                (updateTaskStatus tracer "finished" tId)+  where+    q =+        [sql|WITH pick1 AS (+               SELECT id, payload+               FROM task_queue+               WHERE status IN ('new','suspended')+               ORDER BY priority ASC, created_at DESC+               LIMIT 1+             ), pick2 AS (+               SELECT id, payload+               FROM task_queue+               WHERE status IN ('started') AND age(updated_at) > ?+                 AND NOT (EXISTS (SELECT * FROM pick1))+               ORDER BY priority ASC, updated_at ASC+               LIMIT 1+             )+             SELECT * FROM pick1+               UNION+             SELECT * FROM pick2+             LIMIT 1|]++type Status = Text.Text+type Priority = Text.Text++updateTaskStatus ::+    Tracer IO PGTrace ->+    Status ->+    TaskId ->+    Connection ->+    IO ()+updateTaskStatus tracer st tId conn = do+    let args = (st, priorityIncrement, tId)+    xs :: [(Only Int64)] <- pgQuery (contramap (UpdateStatus st tId) tracer) conn q args+    let xTask = Maybe.listToMaybe xs+    runTracer tracer (UpdatedTask st (fmap fromOnly xTask))+  where+    priorityIncrement :: Int64+    priorityIncrement = case st of "started" -> 20; "suspended" -> 100; _ -> 0+    q =+        [sql|UPDATE task_queue+             SET updated_at = NOW(), status = ?, priority = priority + ?+             WHERE id = ?+             RETURNING id|]++-------------------------------------------------------------------------------++data Step task+    = LookupNextTask+    | ClaimTask (TaskHandle task)+    | WorkingOnTask (TaskHandle task)+instance Show (Step task) where+    show LookupNextTask = "LookupNextTask"+    show (ClaimTask t) = "ClaimTask { taskId = " <> show t.taskId <> " }"+    show (WorkingOnTask t) = "WorkingOnTask { " <> show t.taskId <> " }"++data Trace task+    = StepperTrace (Stepper.Trace (Step task) ())+    | TraceSqlStatement Text.Text+    | PrimitiveSql PGTrace+    deriving (Show)++data RunTaskResult+    = Unstarted+    | Final++data RunTask = RunTask++runTask ::+    forall task.+    (Typeable task, Aeson.FromJSON task, Aeson.ToJSON task) =>+    Tracer IO (Trace task) ->+    IO Connection ->+    (task -> Stepper.ExecFunctions RunTaskResult -> IO ()) ->+    Stepper.BaseStepIO RunTask (Stepper.Delayable RunTaskResult)+runTask tracer mkConnection performTask = \complete _ -> do+    lookupNextTask+        (complete $ Inline Unstarted)+        (claimTask (workOnTask complete))+        (Inline ())+  where+    execution :: (a -> Step task) -> (Stepper.ExecFunctions b -> a -> IO ()) -> StepIO a b+    execution f1 =+        Stepper.defineExecution (contramap StepperTrace tracer) f1 (const ())++    lookupNextTask :: IO () -> StepIO () (TaskHandle task)+    lookupNextTask complete = execution (const LookupNextTask) $ \handle () -> do+        t <- nextTask (contramap PrimitiveSql tracer) =<< mkConnection+        case t of+            Nothing -> print ("could not load task" :: String) >> complete+            Just h -> handle.inline h++    claimTask :: StepIO (TaskHandle task) (TaskHandle task)+    claimTask = execution ClaimTask $ \handle task -> do+        task.markStarted =<< mkConnection+        handle.inline task++    workOnTask :: StepIO (TaskHandle task) RunTaskResult+    workOnTask = execution WorkingOnTask $ \handle task -> do+        let f1 v = mkConnection >>= task.markFinished >> handle.inline v+        let f2 delaySpec v = mkConnection >>= task.markSuspended >> handle.delay delaySpec v+        performTask task.taskDefinition (Stepper.ExecFunctions f1 f2)
+ test/Main.hs view
@@ -0,0 +1,4 @@+module Main (main) where++main :: IO ()+main = putStrLn "Test suite not yet implemented."