sunchao commented on code in PR #58678:
URL: https://github.com/apache/spark/pull/58678#discussion_r4096109886


##########
sql/connect/server/src/main/scala/org/apache/spark/sql/connect/service/SparkConnectService.scala:
##########
@@ -472,6 +479,80 @@ object SparkConnectService extends Logging {
       getClass.getName.stripSuffix("$"))
   }
 
+  /**
+   * Build a Netty SslContext for the Spark Connect gRPC server from the 
`spark.ssl.connect.*`
+   * namespace. Returns None when TLS is disabled.
+   *
+   * Server key material is read as PEM (`certChain` + `privateKey` + optional
+   * `privateKeyPassword`). When `needClientAuth=true` the server additionally 
requires and
+   * verifies a client certificate against a JKS trust store (`trustStore` +
+   * `trustStorePassword`), i.e. mutual TLS. Reloading trust manager
+   * (`trustStoreReloadingEnabled`) is not yet honored on Connect: the server 
logs a warning and
+   * loads the trust store statically. `openSslEnabled=true` is rejected with 
a `SparkException`
+   * at startup rather than silently falling back to the JDK provider, so 
operators do not think
+   * they are running on OpenSSL when they are not. JKS server key material, 
protocol/cipher
+   * overrides, PEM trust anchors, and OpenSSL are follow-ups.
+   *
+   * Matches the `spark.ssl.rpc.enabled` precedent: does NOT inherit 
`spark.ssl.enabled`; must be
+   * opted into explicitly.
+   */
+  private[service] def buildConnectSslContext(sm: SecurityManager): 
Option[SslContext] = {
+    val opts = sm.getSSLOptions("connect")
+    if (!opts.enabled) {
+      logDebug("Spark Connect gRPC server: TLS disabled")
+      return None
+    }
+    if (opts.openSslEnabled) {
+      throw new SparkException(
+        "spark.ssl.connect.openSslEnabled=true is not yet supported " +
+          "on the Spark Connect server; unset it or set it to false to use the 
JDK SSL provider")
+    }
+    val cert = opts.certChain.getOrElse(
+      throw new SparkException(
+        "spark.ssl.connect.enabled=true but spark.ssl.connect.certChain is not 
set"))
+    val key = opts.privateKey.getOrElse(
+      throw new SparkException(
+        "spark.ssl.connect.enabled=true but spark.ssl.connect.privateKey is 
not set"))
+    val builder = SslContextBuilder.forServer(cert, key, 
opts.privateKeyPassword.orNull)
+    val trustManagerFactory = opts.trustStore.map { ts =>
+      try {
+        val ksType = opts.trustStoreType.getOrElse(KeyStore.getDefaultType)
+        val ks = KeyStore.getInstance(ksType)
+        val passwordChars = opts.trustStorePassword.map(_.toCharArray).orNull
+        Utils.tryWithResource(Files.newInputStream(ts.toPath))(ks.load(_, 
passwordChars))
+        val tmf = 
TrustManagerFactory.getInstance(TrustManagerFactory.getDefaultAlgorithm)
+        tmf.init(ks)
+        tmf
+      } catch {
+        case NonFatal(e) =>
+          throw new SparkException(
+            s"Failed to load 
spark.ssl.connect.trustStore='${ts.getAbsolutePath}': " +
+              e.getMessage,
+            e)
+      }
+    }
+    trustManagerFactory.foreach(builder.trustManager)
+    if (opts.trustStoreReloadingEnabled) {
+      logWarning(
+        "spark.ssl.connect.trustStoreReloadingEnabled=true is not yet 
supported; " +
+          "loading the trust store statically at startup")
+    }
+    val clientAuth = if (opts.needClientAuth) {
+      if (trustManagerFactory.isEmpty) {
+        throw new SparkException(
+          "spark.ssl.connect.needClientAuth=true but " +
+            "spark.ssl.connect.trustStore is not set")
+      }
+      builder.clientAuth(ClientAuth.REQUIRE)
+      "REQUIRE"
+    } else {
+      "NONE"
+    }
+    logInfo("Spark Connect gRPC server: TLS enabled")
+    logDebug(s"Spark Connect gRPC server TLS posture: keyMaterial=PEM, 
clientAuth=$clientAuth")
+    Some(GrpcSslContexts.configure(builder).build())

Review Comment:
   **JDK provider selection remains unenforced.** Automatic provider selection 
can use OpenSSL despite `openSslEnabled=false`. The fix remains explicitly 
selecting `SslProvider.JDK`. [Existing 
comment](https://github.com/apache/spark/pull/58678#discussion_r3975700208).



##########
sql/connect/server/src/main/scala/org/apache/spark/sql/connect/service/SparkConnectService.scala:
##########
@@ -437,6 +442,8 @@ object SparkConnectService extends Logging {
       sb.permitKeepAliveWithoutCalls(true)
       sb.addService(sparkConnectService)
 
+      
buildConnectSslContext(SparkEnv.get.securityManager).foreach(sb.sslContext)

Review Comment:
   **Python `foreachBatch` workers lack TLS configuration.** Their return 
connection cannot perform Spark operations against the TLS-enabled server. The 
listener concern applies to the retained server-side callback path; the current 
Python listener API runs callbacks client-side. [Existing 
comment](https://github.com/apache/spark/pull/58678#discussion_r3975700201).



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to