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 +5/−0
- LICENSE +202/−0
- prodapi-pg.cabal +108/−0
- src/Prod/Pg/DatabaseUtils.hs +113/−0
- src/Prod/Pg/TaskQueue.hs +186/−0
- test/Main.hs +4/−0
+ 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."