(ctx context.Context, opts *drivers.ModelExecuteOptions, connector, optionalBucketURL string, optionalAdditionalConfig map[string]any, logger *zap.Logger)
| 597 | } |
| 598 | |
| 599 | func generateSecretSQL(ctx context.Context, opts *drivers.ModelExecuteOptions, connector, optionalBucketURL string, optionalAdditionalConfig map[string]any, logger *zap.Logger) (string, string, string, error) { |
| 600 | handle, release, err := opts.Env.AcquireConnector(ctx, connector) |
| 601 | if err != nil { |
| 602 | return "", "", "", err |
| 603 | } |
| 604 | defer release() |
| 605 | |
| 606 | safeSecretName := safeName(fmt.Sprintf("%s__%s__secret", opts.ModelName, connector)) |
| 607 | dropSecretSQL := fmt.Sprintf("DROP SECRET IF EXISTS %s", safeSecretName) |
| 608 | connectorType := handle.Driver() |
| 609 | |
| 610 | switch connectorType { |
| 611 | case "s3": |
| 612 | conn, ok := handle.(*s3.Connection) |
| 613 | if !ok { |
| 614 | return "", "", "", fmt.Errorf("internal error: expected s3 connector handle") |
| 615 | } |
| 616 | s3Config := conn.ParsedConfig() |
| 617 | err := mapstructure.WeakDecode(optionalAdditionalConfig, s3Config) |
| 618 | if err != nil { |
| 619 | return "", "", "", fmt.Errorf("failed to parse s3 config properties: %w", err) |
| 620 | } |
| 621 | var sb strings.Builder |
| 622 | sb.WriteString("CREATE OR REPLACE TEMPORARY SECRET ") |
| 623 | sb.WriteString(safeSecretName) |
| 624 | sb.WriteString(" (TYPE S3") |
| 625 | // workaround for issue : https://github.com/duckdb/duckdb-python/issues/398#issuecomment-4258043324 |
| 626 | sb.WriteString(", URL_STYLE path") |
| 627 | |
| 628 | if s3Config.AccessKeyID != "" { |
| 629 | fmt.Fprintf(&sb, ", KEY_ID %s, SECRET %s", safeSQLString(s3Config.AccessKeyID), safeSQLString(s3Config.SecretAccessKey)) |
| 630 | } else if s3Config.AllowHostAccess { |
| 631 | sb.WriteString(", PROVIDER CREDENTIAL_CHAIN, VALIDATION 'none'") |
| 632 | } |
| 633 | |
| 634 | if s3Config.SessionToken != "" { |
| 635 | fmt.Fprintf(&sb, ", SESSION_TOKEN %s", safeSQLString(s3Config.SessionToken)) |
| 636 | } |
| 637 | if s3Config.Endpoint != "" { |
| 638 | uri, err := url.Parse(s3Config.Endpoint) |
| 639 | if err == nil && uri.Scheme != "" { // let duckdb raise an error if the endpoint is invalid |
| 640 | // for duckdb the endpoint should not have a scheme |
| 641 | s3Config.Endpoint = strings.TrimPrefix(s3Config.Endpoint, uri.Scheme+"://") |
| 642 | if uri.Scheme == "http" { |
| 643 | sb.WriteString(", USE_SSL false") |
| 644 | } |
| 645 | } |
| 646 | sb.WriteString(", ENDPOINT ") |
| 647 | sb.WriteString(safeSQLString(s3Config.Endpoint)) |
| 648 | } |
| 649 | if s3Config.Region != "" { |
| 650 | sb.WriteString(", REGION ") |
| 651 | sb.WriteString(safeSQLString(s3Config.Region)) |
| 652 | } else if optionalBucketURL != "" { |
| 653 | // DuckDB does not automatically resolve the region as of 1.2.0 so we try to detect and set the region. |
| 654 | uri, err := globutil.ParseBucketURL(optionalBucketURL) |
| 655 | if err != nil { |
| 656 | return "", "", "", fmt.Errorf("failed to parse path %q: %w", optionalBucketURL, err) |
no test coverage detected