From 934d7d44eb7605870f7df7af11feb2de3f720722 Mon Sep 17 00:00:00 2001 From: saksham-datazip Date: Sat, 1 Aug 2026 14:59:23 +0530 Subject: [PATCH 01/11] chore: added s3 flag --- drivers/db2/go.mod | 5 ++ drivers/kafka/go.mod | 7 +- drivers/mongodb/go.mod | 5 ++ drivers/mssql/go.mod | 5 ++ drivers/mssql/main.go | 1 - drivers/mysql/go.mod | 5 ++ drivers/oracle/go.mod | 5 ++ drivers/postgres/go.mod | 5 ++ go.mod | 19 ++++- protocol/clear.go | 9 +- protocol/discover.go | 11 ++- protocol/root.go | 3 + protocol/s3.go | 101 +++++++++++++++++++++++ protocol/sync.go | 11 ++- tests/kafka/go.mod | 5 ++ utils/s3_utils.go | 177 ++++++++++++++++++++++++++++++++++++++++ 16 files changed, 365 insertions(+), 9 deletions(-) create mode 100644 protocol/s3.go create mode 100644 utils/s3_utils.go diff --git a/drivers/db2/go.mod b/drivers/db2/go.mod index ab6ac3135..ac6270e99 100644 --- a/drivers/db2/go.mod +++ b/drivers/db2/go.mod @@ -24,15 +24,20 @@ require ( github.com/apache/thrift v0.23.0 // indirect github.com/aws/aws-sdk-go v1.55.6 // indirect github.com/aws/aws-sdk-go-v2 v1.39.2 // indirect + github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.1 // indirect github.com/aws/aws-sdk-go-v2/config v1.29.17 // indirect github.com/aws/aws-sdk-go-v2/credentials v1.17.70 // indirect github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.16.32 // indirect github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.9 // indirect github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.9 // indirect github.com/aws/aws-sdk-go-v2/internal/ini v1.8.3 // indirect + github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.9 // indirect github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.1 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.0 // indirect github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.9 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.9 // indirect github.com/aws/aws-sdk-go-v2/service/kms v1.41.1 // indirect + github.com/aws/aws-sdk-go-v2/service/s3 v1.88.4 // indirect github.com/aws/aws-sdk-go-v2/service/sso v1.25.5 // indirect github.com/aws/aws-sdk-go-v2/service/ssooidc v1.30.3 // indirect github.com/aws/aws-sdk-go-v2/service/sts v1.34.0 // indirect diff --git a/drivers/kafka/go.mod b/drivers/kafka/go.mod index 576471c30..1508f0bc1 100644 --- a/drivers/kafka/go.mod +++ b/drivers/kafka/go.mod @@ -13,6 +13,7 @@ require ( github.com/linkedin/goavro/v2 v2.15.0 github.com/twmb/franz-go v1.21.1 github.com/twmb/franz-go/pkg/kadm v1.18.0 + github.com/twmb/franz-go/pkg/kmsg v1.13.1 ) require ( @@ -21,15 +22,20 @@ require ( github.com/apache/thrift v0.23.0 // indirect github.com/aws/aws-sdk-go v1.55.6 // indirect github.com/aws/aws-sdk-go-v2 v1.39.2 // indirect + github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.1 // indirect github.com/aws/aws-sdk-go-v2/config v1.29.17 // indirect github.com/aws/aws-sdk-go-v2/credentials v1.17.70 // indirect github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.16.32 // indirect github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.9 // indirect github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.9 // indirect github.com/aws/aws-sdk-go-v2/internal/ini v1.8.3 // indirect + github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.9 // indirect github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.1 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.0 // indirect github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.9 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.9 // indirect github.com/aws/aws-sdk-go-v2/service/kms v1.41.1 // indirect + github.com/aws/aws-sdk-go-v2/service/s3 v1.88.4 // indirect github.com/aws/aws-sdk-go-v2/service/sso v1.25.5 // indirect github.com/aws/aws-sdk-go-v2/service/ssooidc v1.30.3 // indirect github.com/aws/aws-sdk-go-v2/service/sts v1.34.0 // indirect @@ -81,7 +87,6 @@ require ( github.com/subosito/gotenv v1.6.0 // indirect github.com/tklauser/go-sysconf v0.3.16 // indirect github.com/tklauser/numcpus v0.11.0 // indirect - github.com/twmb/franz-go/pkg/kmsg v1.13.1 // indirect github.com/twmb/murmur3 v1.1.8 // indirect github.com/twpayne/go-geom v1.6.1 // indirect github.com/xitongsys/parquet-go v1.6.2 // indirect diff --git a/drivers/mongodb/go.mod b/drivers/mongodb/go.mod index 75d870974..edd42ec36 100644 --- a/drivers/mongodb/go.mod +++ b/drivers/mongodb/go.mod @@ -19,15 +19,20 @@ require ( github.com/apache/thrift v0.23.0 // indirect github.com/aws/aws-sdk-go v1.55.6 // indirect github.com/aws/aws-sdk-go-v2 v1.39.2 // indirect + github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.1 // indirect github.com/aws/aws-sdk-go-v2/config v1.29.17 // indirect github.com/aws/aws-sdk-go-v2/credentials v1.17.70 // indirect github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.16.32 // indirect github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.9 // indirect github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.9 // indirect github.com/aws/aws-sdk-go-v2/internal/ini v1.8.3 // indirect + github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.9 // indirect github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.1 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.0 // indirect github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.9 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.9 // indirect github.com/aws/aws-sdk-go-v2/service/kms v1.41.1 // indirect + github.com/aws/aws-sdk-go-v2/service/s3 v1.88.4 // indirect github.com/aws/aws-sdk-go-v2/service/sso v1.25.5 // indirect github.com/aws/aws-sdk-go-v2/service/ssooidc v1.30.3 // indirect github.com/aws/aws-sdk-go-v2/service/sts v1.34.0 // indirect diff --git a/drivers/mssql/go.mod b/drivers/mssql/go.mod index 0fd545548..b7bd82253 100644 --- a/drivers/mssql/go.mod +++ b/drivers/mssql/go.mod @@ -36,15 +36,20 @@ require ( github.com/apache/thrift v0.23.0 // indirect github.com/aws/aws-sdk-go v1.55.6 // indirect github.com/aws/aws-sdk-go-v2 v1.39.2 // indirect + github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.1 // indirect github.com/aws/aws-sdk-go-v2/config v1.29.17 // indirect github.com/aws/aws-sdk-go-v2/credentials v1.17.70 // indirect github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.16.32 // indirect github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.9 // indirect github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.9 // indirect github.com/aws/aws-sdk-go-v2/internal/ini v1.8.3 // indirect + github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.9 // indirect github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.1 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.0 // indirect github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.9 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.9 // indirect github.com/aws/aws-sdk-go-v2/service/kms v1.41.1 // indirect + github.com/aws/aws-sdk-go-v2/service/s3 v1.88.4 // indirect github.com/aws/aws-sdk-go-v2/service/sso v1.25.5 // indirect github.com/aws/aws-sdk-go-v2/service/ssooidc v1.30.3 // indirect github.com/aws/aws-sdk-go-v2/service/sts v1.34.0 // indirect diff --git a/drivers/mssql/main.go b/drivers/mssql/main.go index 8d80455ae..bc8c26f5c 100644 --- a/drivers/mssql/main.go +++ b/drivers/mssql/main.go @@ -10,4 +10,3 @@ func main() { defer driver.Close() olake.RegisterDriver(driver) } - diff --git a/drivers/mysql/go.mod b/drivers/mysql/go.mod index 0467bc2dd..e967a7696 100644 --- a/drivers/mysql/go.mod +++ b/drivers/mysql/go.mod @@ -20,15 +20,20 @@ require github.com/jmoiron/sqlx v1.4.0 require ( github.com/apache/arrow-go/v18 v18.2.0 // indirect github.com/aws/aws-sdk-go-v2 v1.39.2 // indirect + github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.1 // indirect github.com/aws/aws-sdk-go-v2/config v1.29.17 // indirect github.com/aws/aws-sdk-go-v2/credentials v1.17.70 // indirect github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.16.32 // indirect github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.9 // indirect github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.9 // indirect github.com/aws/aws-sdk-go-v2/internal/ini v1.8.3 // indirect + github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.9 // indirect github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.1 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.0 // indirect github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.9 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.9 // indirect github.com/aws/aws-sdk-go-v2/service/kms v1.41.1 // indirect + github.com/aws/aws-sdk-go-v2/service/s3 v1.88.4 // indirect github.com/aws/aws-sdk-go-v2/service/sso v1.25.5 // indirect github.com/aws/aws-sdk-go-v2/service/ssooidc v1.30.3 // indirect github.com/aws/aws-sdk-go-v2/service/sts v1.34.0 // indirect diff --git a/drivers/oracle/go.mod b/drivers/oracle/go.mod index fa60050f3..f75a83c39 100644 --- a/drivers/oracle/go.mod +++ b/drivers/oracle/go.mod @@ -22,15 +22,20 @@ require ( github.com/apache/thrift v0.23.0 // indirect github.com/aws/aws-sdk-go v1.55.6 // indirect github.com/aws/aws-sdk-go-v2 v1.39.2 // indirect + github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.1 // indirect github.com/aws/aws-sdk-go-v2/config v1.29.17 // indirect github.com/aws/aws-sdk-go-v2/credentials v1.17.70 // indirect github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.16.32 // indirect github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.9 // indirect github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.9 // indirect github.com/aws/aws-sdk-go-v2/internal/ini v1.8.3 // indirect + github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.9 // indirect github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.1 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.0 // indirect github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.9 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.9 // indirect github.com/aws/aws-sdk-go-v2/service/kms v1.41.1 // indirect + github.com/aws/aws-sdk-go-v2/service/s3 v1.88.4 // indirect github.com/aws/aws-sdk-go-v2/service/sso v1.25.5 // indirect github.com/aws/aws-sdk-go-v2/service/ssooidc v1.30.3 // indirect github.com/aws/aws-sdk-go-v2/service/sts v1.34.0 // indirect diff --git a/drivers/postgres/go.mod b/drivers/postgres/go.mod index 978935305..d7e43e9d7 100644 --- a/drivers/postgres/go.mod +++ b/drivers/postgres/go.mod @@ -22,15 +22,20 @@ require ( github.com/apache/thrift v0.23.0 // indirect github.com/aws/aws-sdk-go v1.55.6 // indirect github.com/aws/aws-sdk-go-v2 v1.39.2 // indirect + github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.1 // indirect github.com/aws/aws-sdk-go-v2/config v1.29.17 // indirect github.com/aws/aws-sdk-go-v2/credentials v1.17.70 // indirect github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.16.32 // indirect github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.9 // indirect github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.9 // indirect github.com/aws/aws-sdk-go-v2/internal/ini v1.8.3 // indirect + github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.9 // indirect github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.1 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.0 // indirect github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.9 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.9 // indirect github.com/aws/aws-sdk-go-v2/service/kms v1.41.1 // indirect + github.com/aws/aws-sdk-go-v2/service/s3 v1.88.4 // indirect github.com/aws/aws-sdk-go-v2/service/sso v1.25.5 // indirect github.com/aws/aws-sdk-go-v2/service/ssooidc v1.30.3 // indirect github.com/aws/aws-sdk-go-v2/service/sts v1.34.0 // indirect diff --git a/go.mod b/go.mod index a500a8fe3..de85ea3d1 100644 --- a/go.mod +++ b/go.mod @@ -4,7 +4,9 @@ go 1.25.12 require ( github.com/aws/aws-sdk-go-v2/config v1.29.17 + github.com/aws/aws-sdk-go-v2/credentials v1.17.70 github.com/aws/aws-sdk-go-v2/service/kms v1.41.1 + github.com/aws/aws-sdk-go-v2/service/s3 v1.88.4 github.com/go-mysql-org/go-mysql v1.11.0 github.com/go-playground/locales v0.14.1 github.com/go-playground/universal-translator v0.18.1 @@ -15,6 +17,7 @@ require ( github.com/jackc/pgx/v5 v5.9.2 github.com/jmoiron/sqlx v1.4.0 github.com/linkedin/goavro/v2 v2.15.0 + github.com/minio/minio-go/v7 v7.0.34 github.com/parquet-go/parquet-go v0.29.0 github.com/rs/zerolog v1.34.0 github.com/spf13/cobra v1.9.1 @@ -33,24 +36,37 @@ require ( ) require ( + github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.1 // indirect + github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.9 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.0 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.9 // indirect + github.com/dustin/go-humanize v1.0.1 // indirect github.com/ebitengine/purego v0.10.0 // indirect github.com/go-ole/go-ole v1.2.6 // indirect github.com/golang/snappy v1.0.0 // indirect github.com/jackc/pgx/v4 v4.18.3 // indirect + github.com/json-iterator/go v1.1.12 // indirect github.com/klauspost/asmfmt v1.3.2 // indirect github.com/lufia/plan9stats v0.0.0-20211012122336-39d0f177ccd0 // indirect github.com/minio/asm2plan9s v0.0.0-20200509001527-cdd76441f9d8 // indirect github.com/minio/c2goasm v0.0.0-20190812172519-36a3d3bbc4f3 // indirect + github.com/minio/md5-simd v1.1.2 // indirect + github.com/minio/sha256-simd v1.0.0 // indirect + github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd // indirect + github.com/modern-go/reflect2 v1.0.2 // indirect github.com/parquet-go/bitpack v1.0.0 // indirect github.com/parquet-go/jsonlite v1.0.0 // indirect github.com/power-devops/perfstat v0.0.0-20240221224432-82ca36839d55 // indirect github.com/rogpeppe/go-internal v1.14.1 // indirect + github.com/rs/xid v1.6.0 // indirect + github.com/sirupsen/logrus v1.9.4 // indirect github.com/stretchr/objx v0.5.3 // indirect github.com/tklauser/go-sysconf v0.3.16 // indirect github.com/tklauser/numcpus v0.11.0 // indirect github.com/twpayne/go-geom v1.6.1 // indirect github.com/yusufpapurcu/wmi v1.2.4 // indirect golang.org/x/telemetry v0.0.0-20260625142307-59b4966ccb57 // indirect + gopkg.in/ini.v1 v1.66.6 // indirect ) require ( @@ -63,8 +79,7 @@ require ( github.com/andybalholm/brotli v1.1.1 // indirect github.com/apache/arrow-go/v18 v18.2.0 github.com/apache/thrift v0.23.0 // indirect - github.com/aws/aws-sdk-go-v2 v1.39.2 // indirect - github.com/aws/aws-sdk-go-v2/credentials v1.17.70 // indirect + github.com/aws/aws-sdk-go-v2 v1.39.2 github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.16.32 // indirect github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.9 // indirect github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.9 // indirect diff --git a/protocol/clear.go b/protocol/clear.go index 478c44678..90559942c 100644 --- a/protocol/clear.go +++ b/protocol/clear.go @@ -40,7 +40,14 @@ var clearCmd = &cobra.Command{ } return nil }, - RunE: func(cmd *cobra.Command, _ []string) error { + RunE: func(cmd *cobra.Command, _ []string) (err error) { + defer func() { + if err != nil { + return + } + err = finalizeS3Upload(cmd.Context()) + }() + selectedStreamsMetadata, err := classifyStreams(catalog, nil, state) if err != nil { return fmt.Errorf("failed to get selected streams for clearing: %s", err) diff --git a/protocol/discover.go b/protocol/discover.go index 4e78a8041..133446d94 100644 --- a/protocol/discover.go +++ b/protocol/discover.go @@ -43,12 +43,19 @@ var discoverCmd = &cobra.Command{ return nil }, - RunE: func(cmd *cobra.Command, _ []string) error { + RunE: func(cmd *cobra.Command, _ []string) (err error) { + defer func() { + if err != nil { + return + } + err = finalizeS3Upload(cmd.Context()) + }() + if streamsPath != "" && differencePath != "" { return compareStreams() } - err := connector.Setup(cmd.Context()) + err = connector.Setup(cmd.Context()) if err != nil { return err } diff --git a/protocol/root.go b/protocol/root.go index 5e9cc7f61..92b59fd4d 100644 --- a/protocol/root.go +++ b/protocol/root.go @@ -45,6 +45,9 @@ var RootCmd = &cobra.Command{ Use: "olake", Short: "root command", RunE: func(cmd *cobra.Command, args []string) error { + if err := resolveS3Paths(cmd.Context()); err != nil { + return err + } // set global variables diff --git a/protocol/s3.go b/protocol/s3.go new file mode 100644 index 000000000..2cef15a63 --- /dev/null +++ b/protocol/s3.go @@ -0,0 +1,101 @@ +package protocol + +import ( + "context" + "fmt" + "os" + "path" + "path/filepath" + + "github.com/datazip-inc/olake/constants" + "github.com/datazip-inc/olake/utils" + "github.com/datazip-inc/olake/utils/logger" + "github.com/spf13/viper" +) + +// s3 job folder derived from the first resolved s3:// flag path (e.g. prod/{hash}). +var s3JobBucket string +var s3JobKeyPrefix string + +func resolveS3Paths(ctx context.Context) error { + if err := resolveS3PathFlag(ctx, &configPath); err != nil { + return err + } + if err := resolveS3PathFlag(ctx, &destinationConfigPath); err != nil { + return err + } + if err := resolveS3PathFlag(ctx, &streamsPath); err != nil { + return err + } + if err := resolveS3PathFlag(ctx, &statePath); err != nil { + return err + } + if err := resolveS3PathFlag(ctx, &differencePath); err != nil { + return err + } + + return nil +} + +func resolveS3PathFlag(ctx context.Context, flagPath *string) error { + if *flagPath == "" || *flagPath == "not-set" { + return nil + } + + s3PathMap, err := utils.ResolveS3Path(ctx, *flagPath) + if err != nil { + return err + } + if s3PathMap.IsS3 && s3JobBucket == "" { + s3JobBucket = s3PathMap.Bucket + s3JobKeyPrefix = path.Dir(s3PathMap.Key) + } + *flagPath = s3PathMap.LocalPath + return nil +} + +func finalizeS3Upload(ctx context.Context) error { + if noSave || s3JobBucket == "" { + return nil + } + + configFolder := viper.GetString(constants.ConfigFolder) + if configFolder == "" { + return nil + } + + files := []struct { + local string + name string + }{ + {filepath.Join(configFolder, "streams.json"), "streams.json"}, + {filepath.Join(configFolder, "state.json"), "state.json"}, + {filepath.Join(configFolder, "stats.json"), "stats.json"}, + {viper.GetString(constants.DifferencePath), "difference_streams.json"}, + } + + for _, file := range files { + if file.local == "" { + continue + } + if _, err := os.Stat(file.local); err != nil { + continue + } + + s3Key := path.Join(s3JobKeyPrefix, file.name) + s3URI := fmt.Sprintf("s3://%s/%s", s3JobBucket, s3Key) + err := utils.UploadFileToS3(ctx, utils.S3PathMapping{ + OriginalPath: s3URI, + LocalPath: file.local, + IsS3: true, + Bucket: s3JobBucket, + Key: s3Key, + }) + if err != nil { + return fmt.Errorf("failed to upload config folder artifacts to S3: %s", err) + } + logger.Infof("uploaded %s to %s", file.name, s3URI) + } + + return nil +} diff --git a/protocol/sync.go b/protocol/sync.go index 24673e1f8..7b3a94b7d 100644 --- a/protocol/sync.go +++ b/protocol/sync.go @@ -87,9 +87,16 @@ var syncCmd = &cobra.Command{ logger.Infof("Running sync with state: %s", stateBytes) return nil }, - RunE: func(cmd *cobra.Command, _ []string) error { + RunE: func(cmd *cobra.Command, _ []string) (err error) { + defer func() { + if err != nil { + return + } + err = finalizeS3Upload(cmd.Context()) + }() + // setup conector first - err := connector.Setup(cmd.Context()) + err = connector.Setup(cmd.Context()) if err != nil { return err } diff --git a/tests/kafka/go.mod b/tests/kafka/go.mod index 3c20bcf6f..a058a5e77 100644 --- a/tests/kafka/go.mod +++ b/tests/kafka/go.mod @@ -15,15 +15,20 @@ require ( require ( github.com/andybalholm/brotli v1.1.1 // indirect github.com/aws/aws-sdk-go-v2 v1.39.2 // indirect + github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.1 // indirect github.com/aws/aws-sdk-go-v2/config v1.29.17 // indirect github.com/aws/aws-sdk-go-v2/credentials v1.17.70 // indirect github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.16.32 // indirect github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.9 // indirect github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.9 // indirect github.com/aws/aws-sdk-go-v2/internal/ini v1.8.3 // indirect + github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.9 // indirect github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.1 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.0 // indirect github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.9 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.9 // indirect github.com/aws/aws-sdk-go-v2/service/kms v1.41.1 // indirect + github.com/aws/aws-sdk-go-v2/service/s3 v1.88.4 // indirect github.com/aws/aws-sdk-go-v2/service/sso v1.25.5 // indirect github.com/aws/aws-sdk-go-v2/service/ssooidc v1.30.3 // indirect github.com/aws/aws-sdk-go-v2/service/sts v1.34.0 // indirect diff --git a/utils/s3_utils.go b/utils/s3_utils.go new file mode 100644 index 000000000..81151c031 --- /dev/null +++ b/utils/s3_utils.go @@ -0,0 +1,177 @@ +package utils + +import ( + "context" + "fmt" + "io" + "os" + "path/filepath" + "strings" + + "github.com/aws/aws-sdk-go-v2/aws" + "github.com/aws/aws-sdk-go-v2/config" + "github.com/aws/aws-sdk-go-v2/credentials" + "github.com/aws/aws-sdk-go-v2/service/s3" +) + +const s3URIPrefix = "s3://" + +// S3PathMapping tracks an S3 URI and its downloaded local copy. +type S3PathMapping struct { + OriginalPath string + LocalPath string + IsS3 bool + Bucket string + Key string +} + +// newS3Client creates a new S3 client. +func newS3Client(ctx context.Context) (*s3.Client, error) { + // Load the AWS config. + configOpts := []func(*config.LoadOptions) error{} + + if region := envFirst("OLAKE_S3_REGION", "AWS_REGION"); region != "" { + configOpts = append(configOpts, config.WithRegion(region)) + } + + accessKey := envFirst("OLAKE_S3_ACCESS_KEY_ID", "AWS_ACCESS_KEY_ID") + secretKey := envFirst("OLAKE_S3_SECRET_ACCESS_KEY", "AWS_SECRET_ACCESS_KEY") + if accessKey != "" && secretKey != "" { + sessionToken := envFirst("OLAKE_S3_SESSION_TOKEN", "AWS_SESSION_TOKEN") + configOpts = append(configOpts, config.WithCredentialsProvider( + credentials.NewStaticCredentialsProvider(accessKey, secretKey, sessionToken), + )) + } + + cfg, err := config.LoadDefaultConfig(ctx, configOpts...) + if err != nil { + return nil, fmt.Errorf("failed to load AWS config: %s", err) + } + + opts := []func(*s3.Options){} + if endpoint := envFirst("OLAKE_S3_ENDPOINT", "AWS_ENDPOINT_URL"); endpoint != "" { + opts = append(opts, func(o *s3.Options) { + o.BaseEndpoint = aws.String(endpoint) + o.UsePathStyle = true + }) + } + + return s3.NewFromConfig(cfg, opts...), nil +} + +// envFirst returns the first non-empty environment variable value from the given keys. +func envFirst(keys ...string) string { + for _, key := range keys { + if value := os.Getenv(key); value != "" { + return value + } + } + return "" +} + +// ResolveS3Path downloads an s3:// path to a local temp file. Local paths are returned unchanged. +func ResolveS3Path(ctx context.Context, path string) (S3PathMapping, error) { + if !IsS3Path(path) { + return S3PathMapping{OriginalPath: path, LocalPath: path}, nil + } + + bucket, key, err := ParseS3URI(path) + if err != nil { + return S3PathMapping{}, err + } + + client, err := newS3Client(ctx) + if err != nil { + return S3PathMapping{}, err + } + + resp, err := client.GetObject(ctx, &s3.GetObjectInput{ + Bucket: aws.String(bucket), + Key: aws.String(key), + }) + if err != nil { + return S3PathMapping{}, fmt.Errorf("failed to download %s: %s", path, err) + } + defer resp.Body.Close() + + localPath := localPathForS3URI(bucket, key) + if err := os.MkdirAll(filepath.Dir(localPath), 0o755); err != nil { + return S3PathMapping{}, fmt.Errorf("failed to create temp dir for %s: %s", path, err) + } + + file, err := os.Create(localPath) + if err != nil { + return S3PathMapping{}, fmt.Errorf("failed to create local file for %s: %s", path, err) + } + + if _, err = io.Copy(file, resp.Body); err != nil { + file.Close() + return S3PathMapping{}, fmt.Errorf("failed to write local file for %s: %s", path, err) + } + if err := file.Close(); err != nil { + return S3PathMapping{}, fmt.Errorf("failed to close local file for %s: %s", path, err) + } + + return S3PathMapping{ + OriginalPath: path, + LocalPath: localPath, + IsS3: true, + Bucket: bucket, + Key: key, + }, nil +} + +// ParseS3URI parses an s3:// URI into a bucket and key. +// Caller must ensure uri is an s3:// path. +func ParseS3URI(uri string) (bucket, key string, err error) { + rest := strings.TrimPrefix(uri, s3URIPrefix) + parts := strings.SplitN(rest, "/", 2) + if len(parts) != 2 || parts[0] == "" || parts[1] == "" { + return "", "", fmt.Errorf("invalid s3 uri: %s", uri) + } + + return parts[0], parts[1], nil +} + +// localPathForS3URI returns the local path for an S3 URI. +func localPathForS3URI(bucket, key string) string { + return filepath.Join(os.TempDir(), "olake", "s3", bucket, key) +} + +// IsS3Path returns true if the path is an S3 URI. +func IsS3Path(path string) bool { + return strings.HasPrefix(path, s3URIPrefix) +} + +// UploadFileToS3 uploads the local copy back to S3 when the original path was s3://. +// Caller must ensure LocalPath exists. +func UploadFileToS3(ctx context.Context, s3PathMap S3PathMapping) error { + if !s3PathMap.IsS3 { + return nil + } + + // Create a new S3 client. + client, err := newS3Client(ctx) + if err != nil { + return err + } + + // Open the local file. + file, err := os.Open(s3PathMap.LocalPath) + if err != nil { + return fmt.Errorf("failed to open local file %s: %s", s3PathMap.LocalPath, err) + } + defer file.Close() + + // Upload the file to S3. + _, err = client.PutObject(ctx, &s3.PutObjectInput{ + Bucket: aws.String(s3PathMap.Bucket), + Key: aws.String(s3PathMap.Key), + Body: file, + }) + if err != nil { + return fmt.Errorf("failed to upload %s: %s", s3PathMap.OriginalPath, err) + } + + return nil +} From 744c6f5091d25ddebc5c58e07452abd316d9d062 Mon Sep 17 00:00:00 2001 From: saksham-datazip Date: Sat, 1 Aug 2026 20:03:08 +0530 Subject: [PATCH 02/11] chore: resolved panic --- go.mod | 22 +++++++++++----------- 1 file changed, 11 insertions(+), 11 deletions(-) diff --git a/go.mod b/go.mod index de85ea3d1..e3d05602d 100644 --- a/go.mod +++ b/go.mod @@ -6,7 +6,7 @@ require ( github.com/aws/aws-sdk-go-v2/config v1.29.17 github.com/aws/aws-sdk-go-v2/credentials v1.17.70 github.com/aws/aws-sdk-go-v2/service/kms v1.41.1 - github.com/aws/aws-sdk-go-v2/service/s3 v1.88.4 + github.com/aws/aws-sdk-go-v2/service/s3 v1.97.3 github.com/go-mysql-org/go-mysql v1.11.0 github.com/go-playground/locales v0.14.1 github.com/go-playground/universal-translator v0.18.1 @@ -36,10 +36,10 @@ require ( ) require ( - github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.1 // indirect - github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.9 // indirect - github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.0 // indirect - github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.9 // indirect + github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.8 // indirect + github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.22 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.13 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.21 // indirect github.com/dustin/go-humanize v1.0.1 // indirect github.com/ebitengine/purego v0.10.0 // indirect github.com/go-ole/go-ole v1.2.6 // indirect @@ -79,17 +79,17 @@ require ( github.com/andybalholm/brotli v1.1.1 // indirect github.com/apache/arrow-go/v18 v18.2.0 github.com/apache/thrift v0.23.0 // indirect - github.com/aws/aws-sdk-go-v2 v1.39.2 + github.com/aws/aws-sdk-go-v2 v1.41.5 github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.16.32 // indirect - github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.9 // indirect - github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.9 // indirect + github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.21 // indirect + github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.21 // indirect github.com/aws/aws-sdk-go-v2/internal/ini v1.8.3 // indirect - github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.1 // indirect - github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.9 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.7 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.21 // indirect github.com/aws/aws-sdk-go-v2/service/sso v1.25.5 // indirect github.com/aws/aws-sdk-go-v2/service/ssooidc v1.30.3 // indirect github.com/aws/aws-sdk-go-v2/service/sts v1.34.0 // indirect - github.com/aws/smithy-go v1.23.0 // indirect + github.com/aws/smithy-go v1.24.2 // indirect github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc // indirect github.com/fsnotify/fsnotify v1.8.0 // indirect github.com/gabriel-vasile/mimetype v1.4.8 // indirect From 41cccfceb4db1210f26116e0d91d943e0dba5077 Mon Sep 17 00:00:00 2001 From: saksham-datazip Date: Sun, 2 Aug 2026 19:27:14 +0530 Subject: [PATCH 03/11] chore: changed root file --- drivers/kafka/go.mod | 22 +++++++++++----------- protocol/root.go | 8 +++++++- 2 files changed, 18 insertions(+), 12 deletions(-) diff --git a/drivers/kafka/go.mod b/drivers/kafka/go.mod index 1508f0bc1..fb6259f10 100644 --- a/drivers/kafka/go.mod +++ b/drivers/kafka/go.mod @@ -21,25 +21,25 @@ require ( github.com/apache/arrow-go/v18 v18.2.0 // indirect github.com/apache/thrift v0.23.0 // indirect github.com/aws/aws-sdk-go v1.55.6 // indirect - github.com/aws/aws-sdk-go-v2 v1.39.2 // indirect - github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.1 // indirect + github.com/aws/aws-sdk-go-v2 v1.41.5 // indirect + github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.8 // indirect github.com/aws/aws-sdk-go-v2/config v1.29.17 // indirect github.com/aws/aws-sdk-go-v2/credentials v1.17.70 // indirect github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.16.32 // indirect - github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.9 // indirect - github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.9 // indirect + github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.21 // indirect + github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.21 // indirect github.com/aws/aws-sdk-go-v2/internal/ini v1.8.3 // indirect - github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.9 // indirect - github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.1 // indirect - github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.0 // indirect - github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.9 // indirect - github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.9 // indirect + github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.22 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.7 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.13 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.21 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.21 // indirect github.com/aws/aws-sdk-go-v2/service/kms v1.41.1 // indirect - github.com/aws/aws-sdk-go-v2/service/s3 v1.88.4 // indirect + github.com/aws/aws-sdk-go-v2/service/s3 v1.97.3 // indirect github.com/aws/aws-sdk-go-v2/service/sso v1.25.5 // indirect github.com/aws/aws-sdk-go-v2/service/ssooidc v1.30.3 // indirect github.com/aws/aws-sdk-go-v2/service/sts v1.34.0 // indirect - github.com/aws/smithy-go v1.23.0 // indirect + github.com/aws/smithy-go v1.24.2 // indirect github.com/datazip-inc/olake/lib v0.0.0-00010101000000-000000000000 // indirect github.com/ebitengine/purego v0.10.0 // indirect github.com/fsnotify/fsnotify v1.8.0 // indirect diff --git a/protocol/root.go b/protocol/root.go index 92b59fd4d..da5f4d3b2 100644 --- a/protocol/root.go +++ b/protocol/root.go @@ -44,7 +44,7 @@ var ( var RootCmd = &cobra.Command{ Use: "olake", Short: "root command", - RunE: func(cmd *cobra.Command, args []string) error { + PersistentPreRunE: func(cmd *cobra.Command, _ []string) error { if err := resolveS3Paths(cmd.Context()); err != nil { return err } @@ -74,6 +74,9 @@ var RootCmd = &cobra.Command{ logger.Init() telemetry.Init() + return nil + }, + RunE: func(cmd *cobra.Command, args []string) error { if len(args) == 0 { return cmd.Help() } @@ -134,6 +137,9 @@ func signalAwareRootContext(parent context.Context) context.Context { } func init() { + // Run root PersistentPreRunE (S3 path resolution) before child hooks like sync's. + cobra.EnableTraverseRunHooks = true + // TODO: replace --catalog flag with --streams commands = append(commands, specCmd, checkCmd, discoverCmd, syncCmd, clearCmd) RootCmd.PersistentFlags().StringVarP(&configPath, "config", "", "not-set", "(Required) Config for connector") From c39f5fcba0c903bb541266e59fc8acd08e429740 Mon Sep 17 00:00:00 2001 From: saksham-datazip Date: Sun, 2 Aug 2026 21:45:21 +0530 Subject: [PATCH 04/11] chore: resolved logging issue --- go.mod | 10 ---------- protocol/root.go | 12 ++++++------ 2 files changed, 6 insertions(+), 16 deletions(-) diff --git a/go.mod b/go.mod index e3d05602d..f8429291f 100644 --- a/go.mod +++ b/go.mod @@ -17,7 +17,6 @@ require ( github.com/jackc/pgx/v5 v5.9.2 github.com/jmoiron/sqlx v1.4.0 github.com/linkedin/goavro/v2 v2.15.0 - github.com/minio/minio-go/v7 v7.0.34 github.com/parquet-go/parquet-go v0.29.0 github.com/rs/zerolog v1.34.0 github.com/spf13/cobra v1.9.1 @@ -40,33 +39,24 @@ require ( github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.22 // indirect github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.13 // indirect github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.21 // indirect - github.com/dustin/go-humanize v1.0.1 // indirect github.com/ebitengine/purego v0.10.0 // indirect github.com/go-ole/go-ole v1.2.6 // indirect github.com/golang/snappy v1.0.0 // indirect github.com/jackc/pgx/v4 v4.18.3 // indirect - github.com/json-iterator/go v1.1.12 // indirect github.com/klauspost/asmfmt v1.3.2 // indirect github.com/lufia/plan9stats v0.0.0-20211012122336-39d0f177ccd0 // indirect github.com/minio/asm2plan9s v0.0.0-20200509001527-cdd76441f9d8 // indirect github.com/minio/c2goasm v0.0.0-20190812172519-36a3d3bbc4f3 // indirect - github.com/minio/md5-simd v1.1.2 // indirect - github.com/minio/sha256-simd v1.0.0 // indirect - github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd // indirect - github.com/modern-go/reflect2 v1.0.2 // indirect github.com/parquet-go/bitpack v1.0.0 // indirect github.com/parquet-go/jsonlite v1.0.0 // indirect github.com/power-devops/perfstat v0.0.0-20240221224432-82ca36839d55 // indirect github.com/rogpeppe/go-internal v1.14.1 // indirect - github.com/rs/xid v1.6.0 // indirect - github.com/sirupsen/logrus v1.9.4 // indirect github.com/stretchr/objx v0.5.3 // indirect github.com/tklauser/go-sysconf v0.3.16 // indirect github.com/tklauser/numcpus v0.11.0 // indirect github.com/twpayne/go-geom v1.6.1 // indirect github.com/yusufpapurcu/wmi v1.2.4 // indirect golang.org/x/telemetry v0.0.0-20260625142307-59b4966ccb57 // indirect - gopkg.in/ini.v1 v1.66.6 // indirect ) require ( diff --git a/protocol/root.go b/protocol/root.go index da5f4d3b2..2db24c318 100644 --- a/protocol/root.go +++ b/protocol/root.go @@ -45,17 +45,15 @@ var RootCmd = &cobra.Command{ Use: "olake", Short: "root command", PersistentPreRunE: func(cmd *cobra.Command, _ []string) error { - if err := resolveS3Paths(cmd.Context()); err != nil { - return err - } - + // Resolve now as configPaths are needed by logger.Init(), but the error is handled later because the logger is not initialized yet. + s3Err := resolveS3Paths(cmd.Context()) // set global variables viper.SetDefault(constants.ConfigFolder, os.TempDir()) viper.SetDefault(constants.StatePath, filepath.Join(os.TempDir(), "state.json")) viper.SetDefault(constants.StreamsPath, filepath.Join(os.TempDir(), "streams.json")) viper.SetDefault(constants.DifferencePath, filepath.Join(os.TempDir(), "difference_streams.json")) - if !noSave { + if s3Err == nil && !noSave { configFolder := utils.Ternary(configPath == "not-set", filepath.Dir(destinationConfigPath), filepath.Dir(configPath)).(string) streamsPathEnv := utils.Ternary(streamsPath == "", filepath.Join(configFolder, "streams.json"), streamsPath).(string) differencePathEnv := utils.Ternary(streamsPath != "", filepath.Join(filepath.Dir(streamsPath), "difference_streams.json"), filepath.Join(configFolder, "difference_streams.json")).(string) @@ -74,7 +72,9 @@ var RootCmd = &cobra.Command{ logger.Init() telemetry.Init() - return nil + // Checked last so a resolution failure is reported through the + // now-initialized logger instead of being silently discarded. + return s3Err }, RunE: func(cmd *cobra.Command, args []string) error { if len(args) == 0 { From 814bc9d4ed0373bef8272b5221c26541258dd2b4 Mon Sep 17 00:00:00 2001 From: saksham-datazip Date: Wed, 12 Aug 2026 12:35:52 +0530 Subject: [PATCH 05/11] chore: changed logs mode --- drivers/mysql/go.mod | 22 +++++++++++----------- drivers/postgres/go.mod | 22 +++++++++++----------- drivers/s3/docker-compose.yml | 2 +- protocol/root.go | 4 ++-- utils/logger/logger.go | 23 ++++++++++++++++------- 5 files changed, 41 insertions(+), 32 deletions(-) diff --git a/drivers/mysql/go.mod b/drivers/mysql/go.mod index e7509500d..3410d1d41 100644 --- a/drivers/mysql/go.mod +++ b/drivers/mysql/go.mod @@ -16,25 +16,25 @@ require github.com/jmoiron/sqlx v1.4.0 require ( github.com/apache/arrow-go/v18 v18.2.0 // indirect - github.com/aws/aws-sdk-go-v2 v1.39.2 // indirect - github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.1 // indirect + github.com/aws/aws-sdk-go-v2 v1.41.5 // indirect + github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.8 // indirect github.com/aws/aws-sdk-go-v2/config v1.29.17 // indirect github.com/aws/aws-sdk-go-v2/credentials v1.17.70 // indirect github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.16.32 // indirect - github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.9 // indirect - github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.9 // indirect + github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.21 // indirect + github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.21 // indirect github.com/aws/aws-sdk-go-v2/internal/ini v1.8.3 // indirect - github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.9 // indirect - github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.1 // indirect - github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.0 // indirect - github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.9 // indirect - github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.9 // indirect + github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.22 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.7 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.13 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.21 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.21 // indirect github.com/aws/aws-sdk-go-v2/service/kms v1.41.1 // indirect - github.com/aws/aws-sdk-go-v2/service/s3 v1.88.4 // indirect + github.com/aws/aws-sdk-go-v2/service/s3 v1.97.3 // indirect github.com/aws/aws-sdk-go-v2/service/sso v1.25.5 // indirect github.com/aws/aws-sdk-go-v2/service/ssooidc v1.30.3 // indirect github.com/aws/aws-sdk-go-v2/service/sts v1.34.0 // indirect - github.com/aws/smithy-go v1.23.0 // indirect + github.com/aws/smithy-go v1.24.2 // indirect github.com/ebitengine/purego v0.10.0 // indirect github.com/go-ole/go-ole v1.2.6 // indirect github.com/go-viper/mapstructure/v2 v2.4.0 // indirect diff --git a/drivers/postgres/go.mod b/drivers/postgres/go.mod index 2fc784ad6..ac163d808 100644 --- a/drivers/postgres/go.mod +++ b/drivers/postgres/go.mod @@ -20,25 +20,25 @@ require ( github.com/apache/arrow-go/v18 v18.2.0 // indirect github.com/apache/thrift v0.23.0 // indirect github.com/aws/aws-sdk-go v1.55.6 // indirect - github.com/aws/aws-sdk-go-v2 v1.39.2 // indirect - github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.1 // indirect + github.com/aws/aws-sdk-go-v2 v1.41.5 // indirect + github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.8 // indirect github.com/aws/aws-sdk-go-v2/config v1.29.17 // indirect github.com/aws/aws-sdk-go-v2/credentials v1.17.70 // indirect github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.16.32 // indirect - github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.9 // indirect - github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.9 // indirect + github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.21 // indirect + github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.21 // indirect github.com/aws/aws-sdk-go-v2/internal/ini v1.8.3 // indirect - github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.9 // indirect - github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.1 // indirect - github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.0 // indirect - github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.9 // indirect - github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.9 // indirect + github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.22 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.7 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.13 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.21 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.21 // indirect github.com/aws/aws-sdk-go-v2/service/kms v1.41.1 // indirect - github.com/aws/aws-sdk-go-v2/service/s3 v1.88.4 // indirect + github.com/aws/aws-sdk-go-v2/service/s3 v1.97.3 // indirect github.com/aws/aws-sdk-go-v2/service/sso v1.25.5 // indirect github.com/aws/aws-sdk-go-v2/service/ssooidc v1.30.3 // indirect github.com/aws/aws-sdk-go-v2/service/sts v1.34.0 // indirect - github.com/aws/smithy-go v1.23.0 // indirect + github.com/aws/smithy-go v1.24.2 // indirect github.com/ebitengine/purego v0.10.0 // indirect github.com/fsnotify/fsnotify v1.8.0 // indirect github.com/gabriel-vasile/mimetype v1.4.8 // indirect diff --git a/drivers/s3/docker-compose.yml b/drivers/s3/docker-compose.yml index f24dfad7e..dcc1926da 100644 --- a/drivers/s3/docker-compose.yml +++ b/drivers/s3/docker-compose.yml @@ -16,7 +16,7 @@ services: - 9001:9001 - 9000:9000 volumes: - - ../../destination/iceberg/local-test/data/minio-data:/data + - minio-data:/data command: [ "server", "/data", "--console-address", ":9001" ] mc: diff --git a/protocol/root.go b/protocol/root.go index 2db24c318..342dc4178 100644 --- a/protocol/root.go +++ b/protocol/root.go @@ -68,8 +68,8 @@ var RootCmd = &cobra.Command{ viper.Set(constants.EncryptionKey, encryptionKey) } - // logger uses CONFIG_FOLDER - logger.Init() + // logger uses CONFIG_FOLDER; S3 jobs get JSON stdout matching olake.log + logger.Init(s3JobBucket != "") telemetry.Init() // Checked last so a resolution failure is reported through the diff --git a/utils/logger/logger.go b/utils/logger/logger.go index 7f4e498a0..c23bd67d3 100644 --- a/utils/logger/logger.go +++ b/utils/logger/logger.go @@ -179,7 +179,7 @@ func StatsLogger(ctx context.Context, statsFunc func() (int64, int64, int64, int }() } -func Init() { +func Init(s3Mode bool) { // Set up timestamp for log file names currentTimestamp := time.Now().UTC() timestamp := fmt.Sprintf("%d-%02d-%02d_%02d-%02d-%02d", @@ -187,12 +187,15 @@ func Init() { currentTimestamp.Hour(), currentTimestamp.Minute(), currentTimestamp.Second()) // Configure rotating file logs - rotatingFile := &lumberjack.Logger{ - Filename: fmt.Sprintf("%s/logs/sync_%s/olake.log", viper.GetString(constants.ConfigFolder), timestamp), - MaxSize: 100, // Max size in MB - MaxBackups: 5, // Number of old log files to retain - MaxAge: 30, // Days to retain old log files - Compress: true, + var rotatingFile *lumberjack.Logger + if !s3Mode { + rotatingFile = &lumberjack.Logger{ + Filename: fmt.Sprintf("%s/logs/sync_%s/olake.log", viper.GetString(constants.ConfigFolder), timestamp), + MaxSize: 100, // Max size in MB + MaxBackups: 5, // Number of old log files to retain + MaxAge: 30, // Days to retain old log files + Compress: true, + } } zerolog.TimestampFunc = func() time.Time { return time.Now().UTC() } @@ -236,6 +239,12 @@ func Init() { } } + if s3Mode { + // S3 jobs capture stdout; JSON matches olake.log format. + logger = zerolog.New(os.Stdout).With().Timestamp().Logger() + return + } + // MultiWriter (rotating file + console writer per goroutine) multiWriter := zerolog.MultiLevelWriter(rotatingFile, newConsoleWriter()) From 45637d5a03cf6e27522d3b3de32e921cd4f0aa71 Mon Sep 17 00:00:00 2001 From: saksham-datazip Date: Sun, 16 Aug 2026 01:42:47 +0530 Subject: [PATCH 06/11] chore: added seq log --- utils/logger/logger.go | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/utils/logger/logger.go b/utils/logger/logger.go index c23bd67d3..ff234cda1 100644 --- a/utils/logger/logger.go +++ b/utils/logger/logger.go @@ -14,6 +14,7 @@ import ( "regexp" "strings" "sync" + "sync/atomic" "time" "github.com/datazip-inc/olake/constants" @@ -241,7 +242,10 @@ func Init(s3Mode bool) { if s3Mode { // S3 jobs capture stdout; JSON matches olake.log format. - logger = zerolog.New(os.Stdout).With().Timestamp().Logger() + var seq uint64 + logger = zerolog.New(os.Stdout).Hook(zerolog.HookFunc(func(e *zerolog.Event, _ zerolog.Level, _ string) { + e.Uint64("seq", atomic.AddUint64(&seq, 1)) + })).With().Timestamp().Logger() return } From b64718e01100aeed59221357cc47388701cd5644 Mon Sep 17 00:00:00 2001 From: saksham-datazip Date: Sun, 16 Aug 2026 02:37:22 +0530 Subject: [PATCH 07/11] chore: reverted-docker-compose --- drivers/s3/docker-compose.yml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/drivers/s3/docker-compose.yml b/drivers/s3/docker-compose.yml index dcc1926da..f24dfad7e 100644 --- a/drivers/s3/docker-compose.yml +++ b/drivers/s3/docker-compose.yml @@ -16,7 +16,7 @@ services: - 9001:9001 - 9000:9000 volumes: - - minio-data:/data + - ../../destination/iceberg/local-test/data/minio-data:/data command: [ "server", "/data", "--console-address", ":9001" ] mc: From 9d874ffab953da1bbcdb14132f985ffbced6d656 Mon Sep 17 00:00:00 2001 From: saksham-datazip Date: Mon, 17 Aug 2026 00:52:04 +0530 Subject: [PATCH 08/11] chore: removed envfirst function --- utils/s3_utils.go | 21 +++++---------------- 1 file changed, 5 insertions(+), 16 deletions(-) diff --git a/utils/s3_utils.go b/utils/s3_utils.go index 81151c031..7ea4461a2 100644 --- a/utils/s3_utils.go +++ b/utils/s3_utils.go @@ -30,16 +30,15 @@ func newS3Client(ctx context.Context) (*s3.Client, error) { // Load the AWS config. configOpts := []func(*config.LoadOptions) error{} - if region := envFirst("OLAKE_S3_REGION", "AWS_REGION"); region != "" { + if region := os.Getenv("OLAKE_S3_REGION"); region != "" { configOpts = append(configOpts, config.WithRegion(region)) } - accessKey := envFirst("OLAKE_S3_ACCESS_KEY_ID", "AWS_ACCESS_KEY_ID") - secretKey := envFirst("OLAKE_S3_SECRET_ACCESS_KEY", "AWS_SECRET_ACCESS_KEY") + accessKey := os.Getenv("OLAKE_S3_ACCESS_KEY_ID") + secretKey := os.Getenv("OLAKE_S3_SECRET_ACCESS_KEY") if accessKey != "" && secretKey != "" { - sessionToken := envFirst("OLAKE_S3_SESSION_TOKEN", "AWS_SESSION_TOKEN") configOpts = append(configOpts, config.WithCredentialsProvider( - credentials.NewStaticCredentialsProvider(accessKey, secretKey, sessionToken), + credentials.NewStaticCredentialsProvider(accessKey, secretKey, os.Getenv("OLAKE_S3_SESSION_TOKEN")), )) } @@ -49,7 +48,7 @@ func newS3Client(ctx context.Context) (*s3.Client, error) { } opts := []func(*s3.Options){} - if endpoint := envFirst("OLAKE_S3_ENDPOINT", "AWS_ENDPOINT_URL"); endpoint != "" { + if endpoint := os.Getenv("OLAKE_S3_ENDPOINT"); endpoint != "" { opts = append(opts, func(o *s3.Options) { o.BaseEndpoint = aws.String(endpoint) o.UsePathStyle = true @@ -59,16 +58,6 @@ func newS3Client(ctx context.Context) (*s3.Client, error) { return s3.NewFromConfig(cfg, opts...), nil } -// envFirst returns the first non-empty environment variable value from the given keys. -func envFirst(keys ...string) string { - for _, key := range keys { - if value := os.Getenv(key); value != "" { - return value - } - } - return "" -} - // ResolveS3Path downloads an s3:// path to a local temp file. Local paths are returned unchanged. func ResolveS3Path(ctx context.Context, path string) (S3PathMapping, error) { if !IsS3Path(path) { From 1ad576501ef0453192472bfa4427b293eaf51339 Mon Sep 17 00:00:00 2001 From: saksham-datazip Date: Mon, 31 Aug 2026 14:41:43 +0530 Subject: [PATCH 09/11] chore: Added initOnce --- utils/telemetry/telemetry.go | 77 ++++++++++++++++++++---------------- 1 file changed, 42 insertions(+), 35 deletions(-) diff --git a/utils/telemetry/telemetry.go b/utils/telemetry/telemetry.go index 06415c9ad..9e98ace74 100644 --- a/utils/telemetry/telemetry.go +++ b/utils/telemetry/telemetry.go @@ -64,6 +64,11 @@ var ( disabledOnce sync.Once disabled bool + // initOnce keeps Init from starting a second setup goroutine. PersistentPreRunE + // runs twice per process (protocol init Execute, then RegisterDriver Execute); + // a second close(initDone) panics and kills the CLI with exit 2. + initOnce sync.Once + // initDone closes once Init's background setup finishes. Init does network calls, so an // event sent right after start would otherwise find telemetry nil and be dropped silently. initDone = make(chan struct{}) @@ -94,44 +99,46 @@ type LocationInfo struct { } func Init() { - go func() { - defer close(initDone) - // check for disable - if Disabled() { - return - } - ip := getOutboundIP() - eventProps := loadEventProps() - telemetry = &Telemetry{ - httpClient: &http.Client{Timeout: 5 * time.Second}, - userID: getUserID(eventProps), - service: getService(eventProps), - eventProps: eventProps, - platform: platformInfo{ - OS: runtime.GOOS, - Arch: runtime.GOARCH, - OlakeVersion: version.GetOlakeCLIVersion(), - DeviceCPU: fmt.Sprintf("%d cores", runtime.NumCPU()), - }, - ipAddress: ip, - } + initOnce.Do(func() { + go func() { + defer close(initDone) + // check for disable + if Disabled() { + return + } + ip := getOutboundIP() + eventProps := loadEventProps() + telemetry = &Telemetry{ + httpClient: &http.Client{Timeout: 5 * time.Second}, + userID: getUserID(eventProps), + service: getService(eventProps), + eventProps: eventProps, + platform: platformInfo{ + OS: runtime.GOOS, + Arch: runtime.GOARCH, + OlakeVersion: version.GetOlakeCLIVersion(), + DeviceCPU: fmt.Sprintf("%d cores", runtime.NumCPU()), + }, + ipAddress: ip, + } - if ip != ipNotFoundPlaceholder { - ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) - defer cancel() - loc, err := getLocationFromIP(ctx, ip) - if err == nil { - telemetry.locationInfo = &loc - } else { - logger.Debugf("Failed to fetch location for IP %s: %v", ip, err) - telemetry.locationInfo = &LocationInfo{ - Country: "NA", - Region: "NA", - City: "NA", + if ip != ipNotFoundPlaceholder { + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + loc, err := getLocationFromIP(ctx, ip) + if err == nil { + telemetry.locationInfo = &loc + } else { + logger.Debugf("Failed to fetch location for IP %s: %v", ip, err) + telemetry.locationInfo = &LocationInfo{ + Country: "NA", + Region: "NA", + City: "NA", + } } } - } - }() + }() + }) } // send runs an event in the background, recording it so Flush can wait for it. Panics are From dba12df0e61f7eef7a433df76aefa3a0129012a8 Mon Sep 17 00:00:00 2001 From: saksham-datazip Date: Tue, 1 Sep 2026 02:01:32 +0530 Subject: [PATCH 10/11] chore: changed stats logic --- constants/constants.go | 7 +++ protocol/clear.go | 7 +-- protocol/discover.go | 7 +-- protocol/root.go | 5 +- protocol/s3.go | 101 -------------------------------- protocol/sync.go | 7 +-- utils/logger/logger.go | 8 +++ utils/s3/s3.go | 112 ++++++++++++++++++++++++++++++++++++ utils/s3_utils.go | 127 +++++++++++++++++++---------------------- 9 files changed, 193 insertions(+), 188 deletions(-) delete mode 100644 protocol/s3.go create mode 100644 utils/s3/s3.go diff --git a/constants/constants.go b/constants/constants.go index 05cc3d395..a4a81d2dc 100644 --- a/constants/constants.go +++ b/constants/constants.go @@ -22,6 +22,13 @@ const ( StringifiedData = "data" DefaultReadPreference = "secondaryPreferred" EncryptionKey = "OLAKE_ENCRYPTION_KEY" + EnvStorageMode = "OLAKE_STORAGE_MODE" + StorageModeS3 = "s3" + EnvS3Region = "OLAKE_S3_REGION" + EnvS3AccessKeyID = "OLAKE_S3_ACCESS_KEY_ID" + EnvS3SecretAccessKey = "OLAKE_S3_SECRET_ACCESS_KEY" + EnvS3SessionToken = "OLAKE_S3_SESSION_TOKEN" + EnvS3Endpoint = "OLAKE_S3_ENDPOINT" ConfigFolder = "CONFIG_FOLDER" StatePath = "STATE_PATH" StreamsPath = "STREAMS_PATH" diff --git a/protocol/clear.go b/protocol/clear.go index 2cf48e3a9..1c42f4e68 100644 --- a/protocol/clear.go +++ b/protocol/clear.go @@ -42,12 +42,7 @@ var clearCmd = &cobra.Command{ return nil }, RunE: func(cmd *cobra.Command, _ []string) (err error) { - defer func() { - if err != nil { - return - } - err = finalizeS3Upload(cmd.Context()) - }() + defer utils.FinalizeS3Upload(cmd.Context(), &err, noSave) selectedStreamsMetadata, err := classifyStreams(catalog, nil, state) if err != nil { diff --git a/protocol/discover.go b/protocol/discover.go index 511a29e66..b9cb8b156 100644 --- a/protocol/discover.go +++ b/protocol/discover.go @@ -45,12 +45,7 @@ var discoverCmd = &cobra.Command{ return nil }, RunE: func(cmd *cobra.Command, _ []string) (err error) { - defer func() { - if err != nil { - return - } - err = finalizeS3Upload(cmd.Context()) - }() + defer utils.FinalizeS3Upload(cmd.Context(), &err, noSave) if streamsPath != "" && differencePath != "" { return compareStreams() diff --git a/protocol/root.go b/protocol/root.go index 0ef946778..7fab32471 100644 --- a/protocol/root.go +++ b/protocol/root.go @@ -14,6 +14,7 @@ import ( "github.com/datazip-inc/olake/utils" "github.com/datazip-inc/olake/utils/errs" "github.com/datazip-inc/olake/utils/logger" + "github.com/datazip-inc/olake/utils/s3" "github.com/datazip-inc/olake/utils/telemetry" "github.com/spf13/cobra" "github.com/spf13/viper" @@ -46,7 +47,7 @@ var RootCmd = &cobra.Command{ Short: "root command", PersistentPreRunE: func(cmd *cobra.Command, _ []string) error { // Resolve now as configPaths are needed by logger.Init(), but the error is handled later because the logger is not initialized yet. - s3Err := resolveS3Paths(cmd.Context()) + s3Err := utils.ResolveS3Paths(cmd.Context(), []*string{&configPath, &destinationConfigPath, &streamsPath, &statePath, &differencePath}) // set global variables viper.SetDefault(constants.ConfigFolder, os.TempDir()) viper.SetDefault(constants.StatePath, filepath.Join(os.TempDir(), "state.json")) @@ -68,7 +69,7 @@ var RootCmd = &cobra.Command{ } // logger uses CONFIG_FOLDER; S3 jobs get JSON stdout matching olake.log - logger.Init(s3JobBucket != "") + logger.Init(s3.IsS3Job()) telemetry.Init() // Checked last so a resolution failure is reported through the diff --git a/protocol/s3.go b/protocol/s3.go deleted file mode 100644 index 2cef15a63..000000000 --- a/protocol/s3.go +++ /dev/null @@ -1,101 +0,0 @@ -package protocol - -import ( - "context" - "fmt" - "os" - "path" - "path/filepath" - - "github.com/datazip-inc/olake/constants" - "github.com/datazip-inc/olake/utils" - "github.com/datazip-inc/olake/utils/logger" - "github.com/spf13/viper" -) - -// s3 job folder derived from the first resolved s3:// flag path (e.g. prod/{hash}). -var s3JobBucket string -var s3JobKeyPrefix string - -func resolveS3Paths(ctx context.Context) error { - if err := resolveS3PathFlag(ctx, &configPath); err != nil { - return err - } - if err := resolveS3PathFlag(ctx, &destinationConfigPath); err != nil { - return err - } - if err := resolveS3PathFlag(ctx, &streamsPath); err != nil { - return err - } - if err := resolveS3PathFlag(ctx, &statePath); err != nil { - return err - } - if err := resolveS3PathFlag(ctx, &differencePath); err != nil { - return err - } - - return nil -} - -func resolveS3PathFlag(ctx context.Context, flagPath *string) error { - if *flagPath == "" || *flagPath == "not-set" { - return nil - } - - s3PathMap, err := utils.ResolveS3Path(ctx, *flagPath) - if err != nil { - return err - } - if s3PathMap.IsS3 && s3JobBucket == "" { - s3JobBucket = s3PathMap.Bucket - s3JobKeyPrefix = path.Dir(s3PathMap.Key) - } - *flagPath = s3PathMap.LocalPath - return nil -} - -func finalizeS3Upload(ctx context.Context) error { - if noSave || s3JobBucket == "" { - return nil - } - - configFolder := viper.GetString(constants.ConfigFolder) - if configFolder == "" { - return nil - } - - files := []struct { - local string - name string - }{ - {filepath.Join(configFolder, "streams.json"), "streams.json"}, - {filepath.Join(configFolder, "state.json"), "state.json"}, - {filepath.Join(configFolder, "stats.json"), "stats.json"}, - {viper.GetString(constants.DifferencePath), "difference_streams.json"}, - } - - for _, file := range files { - if file.local == "" { - continue - } - if _, err := os.Stat(file.local); err != nil { - continue - } - - s3Key := path.Join(s3JobKeyPrefix, file.name) - s3URI := fmt.Sprintf("s3://%s/%s", s3JobBucket, s3Key) - err := utils.UploadFileToS3(ctx, utils.S3PathMapping{ - OriginalPath: s3URI, - LocalPath: file.local, - IsS3: true, - Bucket: s3JobBucket, - Key: s3Key, - }) - if err != nil { - return fmt.Errorf("failed to upload config folder artifacts to S3: %s", err) - } - logger.Infof("uploaded %s to %s", file.name, s3URI) - } - - return nil -} diff --git a/protocol/sync.go b/protocol/sync.go index d49a1515d..bdd20724a 100644 --- a/protocol/sync.go +++ b/protocol/sync.go @@ -91,12 +91,7 @@ var syncCmd = &cobra.Command{ return nil }, RunE: func(cmd *cobra.Command, _ []string) (err error) { - defer func() { - if err != nil { - return - } - err = finalizeS3Upload(cmd.Context()) - }() + defer utils.FinalizeS3Upload(cmd.Context(), &err, noSave) // setup conector first err = connector.Setup(cmd.Context()) diff --git a/utils/logger/logger.go b/utils/logger/logger.go index ff234cda1..fbd779558 100644 --- a/utils/logger/logger.go +++ b/utils/logger/logger.go @@ -10,6 +10,7 @@ import ( "net/http/httputil" "os" "os/exec" + "path" "path/filepath" "regexp" "strings" @@ -18,6 +19,7 @@ import ( "time" "github.com/datazip-inc/olake/constants" + "github.com/datazip-inc/olake/utils/s3" "github.com/rs/zerolog" "github.com/shirou/gopsutil/v4/cpu" "github.com/shirou/gopsutil/v4/mem" @@ -175,6 +177,12 @@ func StatsLogger(ctx context.Context, statsFunc func() (int64, int64, int64, int if err := FileLogger(stats, "stats", ".json"); err != nil { Fatalf("failed to write stats in file: %s", err) } + if s3.IsS3Job() { + statsPath := filepath.Join(viper.GetString(constants.ConfigFolder), "stats.json") + if err := s3.UploadFileToS3(ctx, statsPath, s3.JobBucket, path.Join(s3.JobPrefix, "stats.json")); err != nil { + Debugf("failed to upload stats.json: %s", err) + } + } } } }() diff --git a/utils/s3/s3.go b/utils/s3/s3.go new file mode 100644 index 000000000..9923ed103 --- /dev/null +++ b/utils/s3/s3.go @@ -0,0 +1,112 @@ +package s3 + +import ( + "context" + "fmt" + "os" + "path" + + "github.com/aws/aws-sdk-go-v2/aws" + "github.com/aws/aws-sdk-go-v2/config" + "github.com/aws/aws-sdk-go-v2/credentials" + awss3 "github.com/aws/aws-sdk-go-v2/service/s3" + "github.com/datazip-inc/olake/constants" +) + +var ( + s3Client *awss3.Client + JobBucket string + JobPrefix string +) + +// IsS3Job reports whether OLAKE_STORAGE_MODE is s3. +func IsS3Job() bool { + return os.Getenv(constants.EnvStorageMode) == constants.StorageModeS3 +} + +// Init initializes the shared S3 client when storage mode is S3. No-op for NFS. +func Init(ctx context.Context) error { + if !IsS3Job() { + return nil + } + + configOpts := []func(*config.LoadOptions) error{} + if region := os.Getenv(constants.EnvS3Region); region != "" { + configOpts = append(configOpts, config.WithRegion(region)) + } + + accessKey := os.Getenv(constants.EnvS3AccessKeyID) + secretKey := os.Getenv(constants.EnvS3SecretAccessKey) + if accessKey != "" && secretKey != "" { + configOpts = append(configOpts, config.WithCredentialsProvider( + credentials.NewStaticCredentialsProvider(accessKey, secretKey, os.Getenv(constants.EnvS3SessionToken)), + )) + } + + awsCfg, err := config.LoadDefaultConfig(ctx, configOpts...) + if err != nil { + return fmt.Errorf("failed to load AWS config: %s", err) + } + + var s3Opts []func(*awss3.Options) + if endpoint := os.Getenv(constants.EnvS3Endpoint); endpoint != "" { + // Path-style is required for MinIO and other S3-compatible endpoints. + s3Opts = append(s3Opts, func(o *awss3.Options) { + o.BaseEndpoint = aws.String(endpoint) + o.UsePathStyle = true + }) + } + + s3Client = awss3.NewFromConfig(awsCfg, s3Opts...) + return nil +} + +func getS3Client() (*awss3.Client, error) { + if s3Client == nil { + return nil, fmt.Errorf("s3 storage not initialized") + } + return s3Client, nil +} + +// GetObject downloads an object and records the job prefix from the first key. +func GetObject(ctx context.Context, bucket, key string) (*awss3.GetObjectOutput, error) { + if JobBucket == "" { + JobBucket = bucket + JobPrefix = path.Dir(key) + } + + client, err := getS3Client() + if err != nil { + return nil, err + } + return client.GetObject(ctx, &awss3.GetObjectInput{ + Bucket: aws.String(bucket), + Key: aws.String(key), + }) +} + +// UploadFile uploads a local file to bucket/key. +func UploadFileToS3(ctx context.Context, localPath, bucket, key string) error { + client, err := getS3Client() + if err != nil { + return err + } + + file, err := os.Open(localPath) + if err != nil { + return fmt.Errorf("failed to open local file %s: %s", localPath, err) + } + defer file.Close() + + // Upload the file to S3. + _, err = client.PutObject(ctx, &awss3.PutObjectInput{ + Bucket: aws.String(bucket), + Key: aws.String(key), + Body: file, + }) + if err != nil { + return fmt.Errorf("failed to upload s3://%s/%s: %s", bucket, key, err) + } + + return nil +} diff --git a/utils/s3_utils.go b/utils/s3_utils.go index 7ea4461a2..a521d02eb 100644 --- a/utils/s3_utils.go +++ b/utils/s3_utils.go @@ -5,13 +5,14 @@ import ( "fmt" "io" "os" + "path" "path/filepath" "strings" - "github.com/aws/aws-sdk-go-v2/aws" - "github.com/aws/aws-sdk-go-v2/config" - "github.com/aws/aws-sdk-go-v2/credentials" - "github.com/aws/aws-sdk-go-v2/service/s3" + "github.com/datazip-inc/olake/constants" + "github.com/datazip-inc/olake/utils/logger" + s3util "github.com/datazip-inc/olake/utils/s3" + "github.com/spf13/viper" ) const s3URIPrefix = "s3://" @@ -25,39 +26,6 @@ type S3PathMapping struct { Key string } -// newS3Client creates a new S3 client. -func newS3Client(ctx context.Context) (*s3.Client, error) { - // Load the AWS config. - configOpts := []func(*config.LoadOptions) error{} - - if region := os.Getenv("OLAKE_S3_REGION"); region != "" { - configOpts = append(configOpts, config.WithRegion(region)) - } - - accessKey := os.Getenv("OLAKE_S3_ACCESS_KEY_ID") - secretKey := os.Getenv("OLAKE_S3_SECRET_ACCESS_KEY") - if accessKey != "" && secretKey != "" { - configOpts = append(configOpts, config.WithCredentialsProvider( - credentials.NewStaticCredentialsProvider(accessKey, secretKey, os.Getenv("OLAKE_S3_SESSION_TOKEN")), - )) - } - - cfg, err := config.LoadDefaultConfig(ctx, configOpts...) - if err != nil { - return nil, fmt.Errorf("failed to load AWS config: %s", err) - } - - opts := []func(*s3.Options){} - if endpoint := os.Getenv("OLAKE_S3_ENDPOINT"); endpoint != "" { - opts = append(opts, func(o *s3.Options) { - o.BaseEndpoint = aws.String(endpoint) - o.UsePathStyle = true - }) - } - - return s3.NewFromConfig(cfg, opts...), nil -} - // ResolveS3Path downloads an s3:// path to a local temp file. Local paths are returned unchanged. func ResolveS3Path(ctx context.Context, path string) (S3PathMapping, error) { if !IsS3Path(path) { @@ -69,15 +37,7 @@ func ResolveS3Path(ctx context.Context, path string) (S3PathMapping, error) { return S3PathMapping{}, err } - client, err := newS3Client(ctx) - if err != nil { - return S3PathMapping{}, err - } - - resp, err := client.GetObject(ctx, &s3.GetObjectInput{ - Bucket: aws.String(bucket), - Key: aws.String(key), - }) + resp, err := s3util.GetObject(ctx, bucket, key) if err != nil { return S3PathMapping{}, fmt.Errorf("failed to download %s: %s", path, err) } @@ -122,7 +82,6 @@ func ParseS3URI(uri string) (bucket, key string, err error) { return parts[0], parts[1], nil } -// localPathForS3URI returns the local path for an S3 URI. func localPathForS3URI(bucket, key string) string { return filepath.Join(os.TempDir(), "olake", "s3", bucket, key) } @@ -132,35 +91,69 @@ func IsS3Path(path string) bool { return strings.HasPrefix(path, s3URIPrefix) } -// UploadFileToS3 uploads the local copy back to S3 when the original path was s3://. -// Caller must ensure LocalPath exists. -func UploadFileToS3(ctx context.Context, s3PathMap S3PathMapping) error { - if !s3PathMap.IsS3 { +// ResolveS3Paths initializes storage and downloads s3:// flag paths to local temp files. +func ResolveS3Paths(ctx context.Context, flagPaths []*string) error { + if err := s3util.Init(ctx); err != nil { + return err + } + + for _, flagPath := range flagPaths { + if err := resolveS3PathFlag(ctx, flagPath); err != nil { + return err + } + } + return nil +} + +func resolveS3PathFlag(ctx context.Context, flagPath *string) error { + if *flagPath == "" || *flagPath == "not-set" { return nil } - // Create a new S3 client. - client, err := newS3Client(ctx) + s3PathMap, err := ResolveS3Path(ctx, *flagPath) if err != nil { return err } + *flagPath = s3PathMap.LocalPath + return nil +} + +// FinalizeS3Upload uploads local artifacts after a successful run. Deferred with +// a named return so a failed upload is not dropped. Failures skip the upload so +// a partial local write cannot overwrite remote streams/state. +func FinalizeS3Upload(ctx context.Context, err *error, noSave bool) { + if *err != nil || noSave || s3util.JobBucket == "" { + return + } - // Open the local file. - file, err := os.Open(s3PathMap.LocalPath) - if err != nil { - return fmt.Errorf("failed to open local file %s: %s", s3PathMap.LocalPath, err) + statsPath := "" + if configFolder := viper.GetString(constants.ConfigFolder); configFolder != "" { + statsPath = filepath.Join(configFolder, "stats.json") } - defer file.Close() - // Upload the file to S3. - _, err = client.PutObject(ctx, &s3.PutObjectInput{ - Bucket: aws.String(s3PathMap.Bucket), - Key: aws.String(s3PathMap.Key), - Body: file, - }) - if err != nil { - return fmt.Errorf("failed to upload %s: %s", s3PathMap.OriginalPath, err) + files := []struct { + local string + name string + }{ + {viper.GetString(constants.StreamsPath), "streams.json"}, + {viper.GetString(constants.StatePath), "state.json"}, + {statsPath, "stats.json"}, + {viper.GetString(constants.DifferencePath), "difference_streams.json"}, } - return nil + for _, file := range files { + if file.local == "" { + continue + } + if _, statErr := os.Stat(file.local); statErr != nil { + continue + } + + s3Key := path.Join(s3util.JobPrefix, file.name) + if uploadErr := s3util.UploadFileToS3(ctx, file.local, s3util.JobBucket, s3Key); uploadErr != nil { + *err = fmt.Errorf("failed to upload config folder artifacts to S3: %s", uploadErr) + return + } + logger.Infof("uploaded %s to s3://%s/%s", file.name, s3util.JobBucket, s3Key) + } } From d54183c05a0a30964ae4bd7a91fc02888d794a55 Mon Sep 17 00:00:00 2001 From: saksham-datazip Date: Tue, 1 Sep 2026 16:38:09 +0530 Subject: [PATCH 11/11] resolved-merge-conflict --- constants/constants.go | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/constants/constants.go b/constants/constants.go index a4a81d2dc..400e6ddd2 100644 --- a/constants/constants.go +++ b/constants/constants.go @@ -26,8 +26,8 @@ const ( StorageModeS3 = "s3" EnvS3Region = "OLAKE_S3_REGION" EnvS3AccessKeyID = "OLAKE_S3_ACCESS_KEY_ID" - EnvS3SecretAccessKey = "OLAKE_S3_SECRET_ACCESS_KEY" - EnvS3SessionToken = "OLAKE_S3_SESSION_TOKEN" + EnvS3SecretAccessKey = "OLAKE_S3_SECRET_ACCESS_KEY" // #nosec G101 -- env var name, not a credential + EnvS3SessionToken = "OLAKE_S3_SESSION_TOKEN" // #nosec G101 -- env var name, not a credential EnvS3Endpoint = "OLAKE_S3_ENDPOINT" ConfigFolder = "CONFIG_FOLDER" StatePath = "STATE_PATH"