(ctx context.Context, logger zerolog.Logger, specBytes []byte, opts plugin.NewClientOptions)
| 74 | } |
| 75 | |
| 76 | func New(ctx context.Context, logger zerolog.Logger, specBytes []byte, opts plugin.NewClientOptions) (plugin.Client, error) { |
| 77 | c := &Client{ |
| 78 | logger: logger.With().Str("module", "pg-dest").Logger(), |
| 79 | } |
| 80 | if opts.NoConnection { |
| 81 | return c, nil |
| 82 | } |
| 83 | |
| 84 | var s spec.Spec |
| 85 | if err := json.Unmarshal(specBytes, &s); err != nil { |
| 86 | return nil, err |
| 87 | } |
| 88 | s.SetDefaults() |
| 89 | if err := s.Validate(); err != nil { |
| 90 | return nil, err |
| 91 | } |
| 92 | c.spec = &s |
| 93 | c.batchSize = s.BatchSize |
| 94 | c.logger.Info().Str("pgx_log_level", s.PgxLogLevel.String()).Msg("Initializing postgresql destination") |
| 95 | |
| 96 | pgxConfig, err := pgxpool.ParseConfig(s.ConnectionString) |
| 97 | if err != nil { |
| 98 | return nil, fmt.Errorf("failed to parse connection string %w", err) |
| 99 | } |
| 100 | pgxConfig.AfterConnect = func(ctx context.Context, conn *pgx.Conn) error { |
| 101 | return nil |
| 102 | } |
| 103 | |
| 104 | pgxConfig.ConnConfig.Tracer = &tracelog.TraceLog{ |
| 105 | Logger: pgx_zero_log.NewLogger(c.logger), |
| 106 | LogLevel: s.PgxLogLevel.LogLevel(), |
| 107 | } |
| 108 | // maybe expose this to the user? |
| 109 | pgxConfig.ConnConfig.RuntimeParams["timezone"] = "UTC" |
| 110 | if s.HasLakebaseConfig() { |
| 111 | if err := configureLakebase(pgxConfig, s.Lakebase); err != nil { |
| 112 | return nil, fmt.Errorf("failed to configure lakebase: %w", err) |
| 113 | } |
| 114 | } |
| 115 | c.conn, err = pgxpool.NewWithConfig(ctx, pgxConfig) |
| 116 | if err != nil { |
| 117 | return nil, fmt.Errorf("failed to connect to postgresql: %w", err) |
| 118 | } |
| 119 | |
| 120 | c.currentDatabaseName, err = currentDatabase(ctx, c.conn) |
| 121 | if err != nil { |
| 122 | return nil, fmt.Errorf("failed to get current database: %w", err) |
| 123 | } |
| 124 | c.currentSchemaName, err = currentSchema(ctx, c.conn) |
| 125 | if err != nil { |
| 126 | return nil, fmt.Errorf("failed to get current schema: %w", err) |
| 127 | } |
| 128 | c.pgType, err = c.getPgType(ctx) |
| 129 | if err != nil { |
| 130 | return nil, fmt.Errorf("failed to get database type: %w", err) |
| 131 | } |
| 132 | // Initialize embeddings requester if pgvector is configured and supported |
| 133 | if c.hasPgVectorConfig() { |
nothing calls this directly
no test coverage detected