{-# LANGUAGE AllowAmbiguousTypes #-}
{-# LANGUAGE FlexibleContexts #-}
{-# LANGUAGE GADTs #-}
{-# LANGUAGE NumericUnderscores #-}
{-# LANGUAGE OverloadedStrings #-}
{-# LANGUAGE ScopedTypeVariables #-}

module DataFrame.Lazy.Internal.DataFrame where

import qualified Data.Text as T
import DataFrame.IO.CSV (CsvReader, readSeparated)
import qualified DataFrame.Internal.Column as C
import qualified DataFrame.Internal.DataFrame as D
import qualified DataFrame.Internal.Expression as E
import DataFrame.Lazy.Internal.Executor (execute)
import DataFrame.Lazy.Internal.LogicalPlan (
    DataSource (..),
    LogicalPlan (..),
    SortOrder (..),
 )
import qualified DataFrame.Lazy.Internal.Optimizer as Opt
import DataFrame.Operations.Join (JoinType)
import DataFrame.Schema (Schema)

{- | A lazy query that has not been executed yet: a 'LogicalPlan' tree whose
execution is deferred until 'runDataFrame' is called.
-}
data LazyDataFrame = LazyDataFrame
    { LazyDataFrame -> LogicalPlan
plan :: LogicalPlan
    , LazyDataFrame -> Int
batchSize :: Int
    }

instance Show LazyDataFrame where
    show :: LazyDataFrame -> String
show LazyDataFrame
ldf =
        String
"LazyDataFrame { batchSize = "
            String -> ShowS
forall a. Semigroup a => a -> a -> a
<> (Int -> String
forall a. Show a => a -> String
show (LazyDataFrame -> Int
batchSize LazyDataFrame
ldf) String -> ShowS
forall a. Semigroup a => a -> a -> a
<> (String
", plan = " String -> ShowS
forall a. Semigroup a => a -> a -> a
<> (LogicalPlan -> String
forall a. Show a => a -> String
show (LazyDataFrame -> LogicalPlan
plan LazyDataFrame
ldf) String -> ShowS
forall a. Semigroup a => a -> a -> a
<> String
" }")))

-- ---------------------------------------------------------------------------
-- Entry point
-- ---------------------------------------------------------------------------

{- | Execute the lazy query: optimise the logical plan, then stream-execute
the resulting physical plan into a fully-materialised 'D.DataFrame'.
-}
runDataFrame :: LazyDataFrame -> IO D.DataFrame
runDataFrame :: LazyDataFrame -> IO DataFrame
runDataFrame LazyDataFrame
ldf = PhysicalPlan -> IO DataFrame
execute (Int -> LogicalPlan -> PhysicalPlan
Opt.optimize (LazyDataFrame -> Int
batchSize LazyDataFrame
ldf) (LazyDataFrame -> LogicalPlan
plan LazyDataFrame
ldf))

-- ---------------------------------------------------------------------------
-- Builders that construct the logical plan tree
-- ---------------------------------------------------------------------------

-- | Lift an already-loaded eager 'D.DataFrame' into the lazy plan.
fromDataFrame :: D.DataFrame -> LazyDataFrame
fromDataFrame :: DataFrame -> LazyDataFrame
fromDataFrame DataFrame
df = LazyDataFrame{plan :: LogicalPlan
plan = DataFrame -> LogicalPlan
SourceDF DataFrame
df, batchSize :: Int
batchSize = Int
1_000_000}

{- | Scan a CSV file with the default comma separator and the in-tree
strict reader.  For the SIMD reader use 'scanCsvWith'.

The 'Schema' both types and selects: only the columns it names are read,
matching 'scanParquet'.

==== __Example__
@
ghci> schema = D.makeSchema [("id", D.schemaType \@Int), ("name", D.schemaType \@Text)]
ghci> L.runDataFrame (L.scanCsv schema \"customers.csv\")

@
-}
scanCsv :: Schema -> T.Text -> LazyDataFrame
scanCsv :: Schema -> Text -> LazyDataFrame
scanCsv = CsvReader -> Schema -> Text -> LazyDataFrame
scanCsvWith CsvReader
readSeparated

{- | Like 'scanCsv' but with an explicit CSV reader (e.g. the SIMD reader
@fastReadCsvWithOpts@ from @dataframe-fastcsv@). The scan derives the
reader's 'DataFrame.IO.CSV.ReadOptions' from the schema and separator, so
any 'CsvReader' projects.

==== __Example__
@
ghci> import qualified DataFrame.IO.CSV.Fast as Fast
ghci> L.runDataFrame (L.scanCsvWith Fast.fastReadCsvWithOpts schema \"customers.csv\")

@
-}
scanCsvWith :: CsvReader -> Schema -> T.Text -> LazyDataFrame
scanCsvWith :: CsvReader -> Schema -> Text -> LazyDataFrame
scanCsvWith CsvReader
reader Schema
schema Text
path =
    LazyDataFrame
        { plan :: LogicalPlan
plan = DataSource -> Schema -> LogicalPlan
Scan (String -> Char -> CsvReader -> DataSource
CsvSource (Text -> String
T.unpack Text
path) Char
',' CsvReader
reader) Schema
schema
        , batchSize :: Int
batchSize = Int
1_000_000
        }

{- | Like 'scanCsvWith', but the file is read in bounded-memory windows
instead of one pass per chunk — for files too large to hold in memory even
after the schema's projection.

==== __Example__
@
ghci> L.runDataFrame (L.scanCsvStreamingWith Fast.fastReadCsvWithOpts schema \"huge.csv\")

@
-}
scanCsvStreamingWith :: CsvReader -> Schema -> T.Text -> LazyDataFrame
scanCsvStreamingWith :: CsvReader -> Schema -> Text -> LazyDataFrame
scanCsvStreamingWith CsvReader
reader Schema
schema Text
path =
    LazyDataFrame
        { plan :: LogicalPlan
plan = DataSource -> Schema -> LogicalPlan
Scan (String -> Char -> CsvReader -> DataSource
CsvSourceStreaming (Text -> String
T.unpack Text
path) Char
',' CsvReader
reader) Schema
schema
        , batchSize :: Int
batchSize = Int
1_000_000
        }

{- | Scan a character-separated file with the default strict reader.

==== __Example__
@
ghci> L.runDataFrame (L.scanSeparated ';' schema \"customers.txt\")

@
-}
scanSeparated :: Char -> Schema -> T.Text -> LazyDataFrame
scanSeparated :: Char -> Schema -> Text -> LazyDataFrame
scanSeparated = CsvReader -> Char -> Schema -> Text -> LazyDataFrame
scanSeparatedWith CsvReader
readSeparated

{- | Like 'scanSeparated' but with an explicit CSV reader.

==== __Example__
@
ghci> L.runDataFrame (L.scanSeparatedWith Fast.fastReadCsvWithOpts ';' schema \"customers.txt\")

@
-}
scanSeparatedWith ::
    CsvReader -> Char -> Schema -> T.Text -> LazyDataFrame
scanSeparatedWith :: CsvReader -> Char -> Schema -> Text -> LazyDataFrame
scanSeparatedWith CsvReader
reader Char
sep Schema
schema Text
path =
    LazyDataFrame
        { plan :: LogicalPlan
plan = DataSource -> Schema -> LogicalPlan
Scan (String -> Char -> CsvReader -> DataSource
CsvSource (Text -> String
T.unpack Text
path) Char
sep CsvReader
reader) Schema
schema
        , batchSize :: Int
batchSize = Int
1_000_000
        }

-- | Scan a Parquet file, directory of files, or glob pattern.
scanParquet :: Schema -> T.Text -> LazyDataFrame
scanParquet :: Schema -> Text -> LazyDataFrame
scanParquet Schema
schema Text
path =
    LazyDataFrame
        { plan :: LogicalPlan
plan = DataSource -> Schema -> LogicalPlan
Scan (String -> DataSource
ParquetSource (Text -> String
T.unpack Text
path)) Schema
schema
        , batchSize :: Int
batchSize = Int
1_000_000
        }

-- | Add a computed column (or overwrite an existing one).
derive ::
    (C.Columnable a) => T.Text -> E.Expr a -> LazyDataFrame -> LazyDataFrame
derive :: forall a.
Columnable a =>
Text -> Expr a -> LazyDataFrame -> LazyDataFrame
derive Text
name Expr a
expr LazyDataFrame
ldf =
    LazyDataFrame
ldf{plan = Derive name (E.UExpr expr) (plan ldf)}

-- | Retain only the listed columns.
select :: [T.Text] -> LazyDataFrame -> LazyDataFrame
select :: [Text] -> LazyDataFrame -> LazyDataFrame
select [Text]
cols LazyDataFrame
ldf = LazyDataFrame
ldf{plan = Project cols (plan ldf)}

-- | Keep rows that satisfy the predicate.
filter :: E.Expr Bool -> LazyDataFrame -> LazyDataFrame
filter :: Expr Bool -> LazyDataFrame -> LazyDataFrame
filter Expr Bool
cond LazyDataFrame
ldf = LazyDataFrame
ldf{plan = Filter cond (plan ldf)}

-- | Join two lazy queries on the given key columns.
join ::
    JoinType ->
    -- | Left join key column name
    T.Text ->
    -- | Right join key column name
    T.Text ->
    -- | Left sub-query
    LazyDataFrame ->
    -- | Right sub-query
    LazyDataFrame ->
    LazyDataFrame
join :: JoinType
-> Text -> Text -> LazyDataFrame -> LazyDataFrame -> LazyDataFrame
join JoinType
jt Text
leftKey Text
rightKey LazyDataFrame
left LazyDataFrame
right =
    LazyDataFrame
        { plan :: LogicalPlan
plan = JoinType
-> Text -> Text -> LogicalPlan -> LogicalPlan -> LogicalPlan
Join JoinType
jt Text
leftKey Text
rightKey (LazyDataFrame -> LogicalPlan
plan LazyDataFrame
left) (LazyDataFrame -> LogicalPlan
plan LazyDataFrame
right)
        , batchSize :: Int
batchSize = LazyDataFrame -> Int
batchSize LazyDataFrame
left
        }

{- | Group by a set of columns and compute aggregate expressions.

Each aggregate expression should use an 'Agg' node (e.g. @sumOf@, @meanOf@).
-}
groupBy ::
    -- | Group-by key columns
    [T.Text] ->
    -- | @[(outputName, aggregateExpr)]@
    [(T.Text, E.UExpr)] ->
    LazyDataFrame ->
    LazyDataFrame
groupBy :: [Text] -> [(Text, UExpr)] -> LazyDataFrame -> LazyDataFrame
groupBy [Text]
keys [(Text, UExpr)]
aggs LazyDataFrame
ldf = LazyDataFrame
ldf{plan = Aggregate keys aggs (plan ldf)}

-- | Sort the result by the given @(column, direction)@ pairs.
sortBy :: [(T.Text, SortOrder)] -> LazyDataFrame -> LazyDataFrame
sortBy :: [(Text, SortOrder)] -> LazyDataFrame -> LazyDataFrame
sortBy [(Text, SortOrder)]
cols LazyDataFrame
ldf = LazyDataFrame
ldf{plan = Sort cols (plan ldf)}

-- | Retain at most @n@ rows.
take :: Int -> LazyDataFrame -> LazyDataFrame
take :: Int -> LazyDataFrame -> LazyDataFrame
take Int
n LazyDataFrame
ldf = LazyDataFrame
ldf{plan = Limit n (plan ldf)}