2021-09-24 01:56:37 +03:00
|
|
|
{-# LANGUAGE DeriveAnyClass #-}
|
2022-03-16 03:39:21 +03:00
|
|
|
{-# LANGUAGE TemplateHaskell #-}
|
2021-04-12 13:18:29 +03:00
|
|
|
{-# LANGUAGE NoGeneralisedNewtypeDeriving #-}
|
|
|
|
{-# OPTIONS_GHC -fno-warn-orphans #-}
|
|
|
|
|
2021-11-04 19:08:33 +03:00
|
|
|
module Hasura.Backends.BigQuery.Source
|
|
|
|
( BigQueryConnSourceConfig (..),
|
2022-02-14 12:45:46 +03:00
|
|
|
RetryOptions (..),
|
2022-02-09 18:26:14 +03:00
|
|
|
BigQueryConnection (..),
|
2021-11-04 19:08:33 +03:00
|
|
|
BigQuerySourceConfig (..),
|
|
|
|
ConfigurationInput (..),
|
|
|
|
ConfigurationInputs (..),
|
|
|
|
ConfigurationJSON (..),
|
|
|
|
GoogleAccessToken (GoogleAccessToken),
|
|
|
|
PKey (unPKey),
|
|
|
|
ServiceAccount (..),
|
|
|
|
TokenResp (..),
|
|
|
|
)
|
|
|
|
where
|
2021-04-12 13:18:29 +03:00
|
|
|
|
2021-09-24 01:56:37 +03:00
|
|
|
import Control.Concurrent.MVar
|
|
|
|
import Crypto.PubKey.RSA.Types qualified as Cry
|
|
|
|
import Data.Aeson qualified as J
|
|
|
|
import Data.Aeson.Casing qualified as J
|
2022-06-08 18:31:28 +03:00
|
|
|
import Data.Aeson.KeyMap qualified as KM
|
2021-09-24 01:56:37 +03:00
|
|
|
import Data.Aeson.TH qualified as J
|
|
|
|
import Data.ByteString.Lazy qualified as BL
|
2022-01-17 13:01:25 +03:00
|
|
|
import Data.Int qualified as Int
|
2021-09-24 01:56:37 +03:00
|
|
|
import Data.Text.Encoding qualified as TE
|
|
|
|
import Data.X509 qualified as X509
|
|
|
|
import Data.X509.Memory qualified as X509
|
|
|
|
import Hasura.Incremental (Cacheable (..))
|
|
|
|
import Hasura.Prelude
|
2021-04-12 13:18:29 +03:00
|
|
|
|
|
|
|
data PKey = PKey
|
2021-09-24 01:56:37 +03:00
|
|
|
{ unPKey :: Cry.PrivateKey,
|
|
|
|
originalBS :: Text
|
2021-04-12 13:18:29 +03:00
|
|
|
}
|
|
|
|
deriving (Show, Eq, Data, Generic, NFData, Hashable)
|
2021-09-24 01:56:37 +03:00
|
|
|
|
2021-04-12 13:18:29 +03:00
|
|
|
deriving instance Generic Cry.PrivateKey -- orphan
|
2021-09-24 01:56:37 +03:00
|
|
|
|
2021-04-12 13:18:29 +03:00
|
|
|
deriving instance Generic Cry.PublicKey -- orphan
|
2021-09-24 01:56:37 +03:00
|
|
|
|
2021-04-12 13:18:29 +03:00
|
|
|
deriving instance J.ToJSON Cry.PrivateKey -- orphan
|
2021-09-24 01:56:37 +03:00
|
|
|
|
2021-04-12 13:18:29 +03:00
|
|
|
deriving instance J.ToJSON Cry.PublicKey -- orphan
|
2021-09-24 01:56:37 +03:00
|
|
|
|
2021-04-12 13:18:29 +03:00
|
|
|
deriving instance Hashable Cry.PrivateKey -- orphan
|
2021-09-24 01:56:37 +03:00
|
|
|
|
2021-04-12 13:18:29 +03:00
|
|
|
deriving instance Hashable Cry.PublicKey -- orphan
|
2021-09-24 01:56:37 +03:00
|
|
|
|
2021-04-12 13:18:29 +03:00
|
|
|
instance J.FromJSON PKey where
|
|
|
|
parseJSON = J.withText "private_key" $ \k ->
|
|
|
|
case X509.readKeyFileFromMemory $ TE.encodeUtf8 k of
|
|
|
|
[X509.PrivKeyRSA k'] -> return $ PKey k' k
|
2021-09-24 01:56:37 +03:00
|
|
|
_ -> fail "unable to parse private key"
|
2021-04-12 13:18:29 +03:00
|
|
|
|
2021-09-24 01:56:37 +03:00
|
|
|
instance J.ToJSON PKey where
|
|
|
|
toJSON PKey {..} = J.toJSON originalBS
|
2021-04-12 13:18:29 +03:00
|
|
|
|
|
|
|
newtype GoogleAccessToken
|
|
|
|
= GoogleAccessToken Text
|
|
|
|
deriving (Show, Eq, J.FromJSON, J.ToJSON, Hashable, Generic, Data, NFData)
|
|
|
|
|
2021-09-24 01:56:37 +03:00
|
|
|
data TokenResp = TokenResp
|
2022-07-29 17:05:03 +03:00
|
|
|
{ _trAccessToken :: GoogleAccessToken,
|
|
|
|
_trExpiresAt :: Integer -- Number of seconds until expiry from `now`, but we add `now` seconds to this for easy tracking
|
2021-09-24 01:56:37 +03:00
|
|
|
}
|
|
|
|
deriving (Eq, Show, Data, NFData, Generic, Hashable)
|
2021-04-12 13:18:29 +03:00
|
|
|
|
|
|
|
instance J.FromJSON TokenResp where
|
2021-09-24 01:56:37 +03:00
|
|
|
parseJSON = J.withObject "TokenResp" $ \o ->
|
|
|
|
TokenResp
|
|
|
|
<$> o J..: "access_token"
|
|
|
|
<*> o J..: "expires_in"
|
2021-04-12 13:18:29 +03:00
|
|
|
|
2021-09-24 01:56:37 +03:00
|
|
|
data ServiceAccount = ServiceAccount
|
2022-07-29 17:05:03 +03:00
|
|
|
{ _saClientEmail :: Text,
|
|
|
|
_saPrivateKey :: PKey,
|
|
|
|
_saProjectId :: Text
|
2021-09-24 01:56:37 +03:00
|
|
|
}
|
|
|
|
deriving (Eq, Show, Data, NFData, Generic, Hashable)
|
2021-04-12 13:18:29 +03:00
|
|
|
|
2021-09-24 01:56:37 +03:00
|
|
|
$(J.deriveJSON (J.aesonDrop 3 J.snakeCase) {J.omitNothingFields = False} ''ServiceAccount)
|
2021-04-12 13:18:29 +03:00
|
|
|
|
|
|
|
data ConfigurationJSON a
|
|
|
|
= FromEnvJSON Text
|
|
|
|
| FromYamlJSON a
|
|
|
|
deriving stock (Show, Eq, Generic)
|
|
|
|
deriving (NFData, Hashable)
|
2021-09-24 01:56:37 +03:00
|
|
|
|
2021-04-12 13:18:29 +03:00
|
|
|
instance J.FromJSON a => J.FromJSON (ConfigurationJSON a) where
|
|
|
|
parseJSON = \case
|
2022-06-08 18:31:28 +03:00
|
|
|
J.Object o | Just (J.String text) <- KM.lookup "from_env" o -> pure (FromEnvJSON text)
|
2021-04-12 13:18:29 +03:00
|
|
|
J.String s -> case J.eitherDecode . BL.fromStrict . TE.encodeUtf8 $ s of
|
2021-09-24 01:56:37 +03:00
|
|
|
Left {} -> fail "error parsing configuration json"
|
2021-04-12 13:18:29 +03:00
|
|
|
Right sa -> pure sa
|
|
|
|
j -> fmap FromYamlJSON (J.parseJSON j)
|
2021-09-24 01:56:37 +03:00
|
|
|
|
2021-04-12 13:18:29 +03:00
|
|
|
instance J.ToJSON a => J.ToJSON (ConfigurationJSON a) where
|
|
|
|
toJSON = \case
|
2021-09-24 01:56:37 +03:00
|
|
|
FromEnvJSON i -> J.object ["from_env" J..= i]
|
2021-04-12 13:18:29 +03:00
|
|
|
FromYamlJSON j -> J.toJSON j
|
|
|
|
|
|
|
|
-- | Configuration inputs when they are a YAML array or an Env var whos value is
|
|
|
|
-- a comma-separated string
|
|
|
|
data ConfigurationInputs
|
2022-07-29 17:05:03 +03:00
|
|
|
= FromYamls [Text]
|
|
|
|
| FromEnvs Text
|
2021-04-12 13:18:29 +03:00
|
|
|
deriving stock (Show, Eq, Generic)
|
|
|
|
deriving (NFData, Hashable)
|
2021-09-24 01:56:37 +03:00
|
|
|
|
2021-04-12 13:18:29 +03:00
|
|
|
instance J.ToJSON ConfigurationInputs where
|
|
|
|
toJSON = \case
|
|
|
|
FromYamls i -> J.toJSON i
|
2021-09-24 01:56:37 +03:00
|
|
|
FromEnvs i -> J.object ["from_env" J..= i]
|
|
|
|
|
2021-04-12 13:18:29 +03:00
|
|
|
instance J.FromJSON ConfigurationInputs where
|
|
|
|
parseJSON = \case
|
2021-09-24 01:56:37 +03:00
|
|
|
J.Object o -> FromEnvs <$> o J..: "from_env"
|
2021-04-12 13:18:29 +03:00
|
|
|
s@(J.Array _) -> FromYamls <$> J.parseJSON s
|
2021-09-24 01:56:37 +03:00
|
|
|
_ -> fail "one of array or object must be provided"
|
2021-04-12 13:18:29 +03:00
|
|
|
|
|
|
|
-- | Configuration input when the YAML value as well as the Env var have
|
|
|
|
-- singlular values
|
|
|
|
data ConfigurationInput
|
2022-07-29 17:05:03 +03:00
|
|
|
= FromYaml Text
|
|
|
|
| FromEnv Text
|
2021-04-12 13:18:29 +03:00
|
|
|
deriving stock (Show, Eq, Generic)
|
|
|
|
deriving (NFData, Hashable)
|
2021-09-24 01:56:37 +03:00
|
|
|
|
2021-04-12 13:18:29 +03:00
|
|
|
instance J.ToJSON ConfigurationInput where
|
|
|
|
toJSON = \case
|
|
|
|
FromYaml i -> J.toJSON i
|
2021-09-24 01:56:37 +03:00
|
|
|
FromEnv i -> J.object ["from_env" J..= i]
|
|
|
|
|
2021-04-12 13:18:29 +03:00
|
|
|
instance J.FromJSON ConfigurationInput where
|
|
|
|
parseJSON = \case
|
2021-09-24 01:56:37 +03:00
|
|
|
J.Object o -> FromEnv <$> o J..: "from_env"
|
2021-04-12 13:18:29 +03:00
|
|
|
s@(J.String _) -> FromYaml <$> J.parseJSON s
|
2021-09-24 01:56:37 +03:00
|
|
|
(J.Number n) -> FromYaml <$> J.parseJSON (J.String (tshow n))
|
|
|
|
_ -> fail "one of string or number or object must be provided"
|
|
|
|
|
|
|
|
data BigQueryConnSourceConfig = BigQueryConnSourceConfig
|
2022-02-14 12:45:46 +03:00
|
|
|
{ _cscServiceAccount :: ConfigurationJSON ServiceAccount,
|
|
|
|
_cscDatasets :: ConfigurationInputs,
|
|
|
|
_cscProjectId :: ConfigurationInput, -- this is part of service-account.json, but we put it here on purpose
|
|
|
|
_cscGlobalSelectLimit :: Maybe ConfigurationInput,
|
|
|
|
_cscRetryBaseDelay :: Maybe ConfigurationInput,
|
|
|
|
_cscRetryLimit :: Maybe ConfigurationInput
|
2021-09-24 01:56:37 +03:00
|
|
|
}
|
|
|
|
deriving (Eq, Generic, NFData)
|
|
|
|
|
|
|
|
$(J.deriveJSON (J.aesonDrop 4 J.snakeCase) {J.omitNothingFields = True} ''BigQueryConnSourceConfig)
|
|
|
|
|
2021-04-12 13:18:29 +03:00
|
|
|
deriving instance Show BigQueryConnSourceConfig
|
2021-09-24 01:56:37 +03:00
|
|
|
|
2021-04-12 13:18:29 +03:00
|
|
|
deriving instance Hashable BigQueryConnSourceConfig
|
2021-09-24 01:56:37 +03:00
|
|
|
|
2021-04-12 13:18:29 +03:00
|
|
|
instance Cacheable BigQueryConnSourceConfig where
|
|
|
|
unchanged _ = (==)
|
|
|
|
|
2022-02-14 12:45:46 +03:00
|
|
|
data RetryOptions = RetryOptions
|
|
|
|
{ _retryBaseDelay :: Microseconds,
|
|
|
|
_retryNumRetries :: Int
|
|
|
|
}
|
|
|
|
deriving (Eq)
|
|
|
|
|
2022-02-09 18:26:14 +03:00
|
|
|
data BigQueryConnection = BigQueryConnection
|
2022-02-14 12:45:46 +03:00
|
|
|
{ _bqServiceAccount :: ServiceAccount,
|
|
|
|
_bqProjectId :: Text, -- this is part of service-account.json, but we put it here on purpose
|
|
|
|
_bqRetryOptions :: Maybe RetryOptions,
|
|
|
|
_bqAccessTokenMVar :: MVar (Maybe TokenResp)
|
2022-02-09 18:26:14 +03:00
|
|
|
}
|
|
|
|
deriving (Eq)
|
|
|
|
|
2021-09-24 01:56:37 +03:00
|
|
|
data BigQuerySourceConfig = BigQuerySourceConfig
|
2022-02-14 12:45:46 +03:00
|
|
|
{ _scConnection :: BigQueryConnection,
|
|
|
|
_scDatasets :: [Text],
|
|
|
|
_scGlobalSelectLimit :: Int.Int64
|
2021-09-24 01:56:37 +03:00
|
|
|
}
|
|
|
|
deriving (Eq)
|
2021-04-12 13:18:29 +03:00
|
|
|
|
|
|
|
instance Cacheable BigQuerySourceConfig where
|
|
|
|
unchanged _ = (==)
|
2021-09-24 01:56:37 +03:00
|
|
|
|
2021-07-30 18:42:36 +03:00
|
|
|
instance J.ToJSON BigQuerySourceConfig where
|
2021-09-24 01:56:37 +03:00
|
|
|
toJSON BigQuerySourceConfig {..} =
|
2022-02-14 12:45:46 +03:00
|
|
|
J.object $
|
2022-02-09 18:26:14 +03:00
|
|
|
[ "service_account" J..= _bqServiceAccount _scConnection,
|
2021-09-24 01:56:37 +03:00
|
|
|
"datasets" J..= _scDatasets,
|
2022-02-09 18:26:14 +03:00
|
|
|
"project_id" J..= _bqProjectId _scConnection,
|
2021-09-24 01:56:37 +03:00
|
|
|
"global_select_limit" J..= _scGlobalSelectLimit
|
2021-07-30 18:42:36 +03:00
|
|
|
]
|
2022-02-14 12:45:46 +03:00
|
|
|
<> case _bqRetryOptions _scConnection of
|
|
|
|
Just RetryOptions {..} ->
|
|
|
|
[ "base_delay" J..= diffTimeToMicroSeconds (microseconds _retryBaseDelay),
|
|
|
|
"retry_limit" J..= _retryNumRetries
|
|
|
|
]
|
|
|
|
Nothing -> []
|