Skip to content

Commit c1c4b4d

Browse files
committed
Permit using RBAC and Cert Authenticator
1 parent 3b45e72 commit c1c4b4d

1 file changed

Lines changed: 27 additions & 8 deletions

File tree

core/src/main/scala/com/sandinh/couchbase/CBCluster.scala

Lines changed: 27 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@ import java.lang
44
import javax.inject._
55

66
import com.couchbase.client.java.CouchbaseAsyncCluster
7+
import com.couchbase.client.java.auth.{CertAuthenticator, PasswordAuthenticator}
78
import com.couchbase.client.java.document.Document
89
import com.couchbase.client.java.env.CouchbaseEnvironment
910
import com.couchbase.client.java.transcoder.Transcoder
@@ -12,7 +13,6 @@ import com.typesafe.config.Config
1213

1314
import scala.jdk.CollectionConverters._
1415
import scala.concurrent.{Await, Future}
15-
import scala.util.Try
1616
import com.sandinh.couchbase.JavaConverters._
1717
import com.sandinh.rx.Implicits._
1818

@@ -23,11 +23,21 @@ import scala.concurrent.duration._
2323
class CBCluster @Inject() (config: Config) {
2424
val env: CouchbaseEnvironment = CbEnvBuilder(config)
2525

26-
val asJava: CouchbaseAsyncCluster =
27-
CouchbaseAsyncCluster.fromConnectionString(
26+
val asJava: CouchbaseAsyncCluster = {
27+
val cfg = config.getConfig("com.sandinh.couchbase")
28+
val cluster = CouchbaseAsyncCluster.fromConnectionString(
2829
env,
29-
config.getString("com.sandinh.couchbase.connectionString")
30+
cfg.getString("connectionString")
3031
)
32+
if (!cfg.hasPath("user")) cluster
33+
else
34+
cluster.authenticate(
35+
new PasswordAuthenticator(
36+
cfg.getString("user"),
37+
cfg.getString("password")
38+
)
39+
)
40+
}
3141

3242
/** Open bucket with typesafe config load from key com.sandinh.couchbase.buckets.`bucket`
3343
* @param bucket use as a subkey of typesafe config for open bucket.
@@ -46,14 +56,23 @@ class CBCluster @Inject() (config: Config) {
4656
legacyEncodeString: Boolean,
4757
transcoders: Transcoder[_ <: Document[_], _]*
4858
): Future[ScalaBucket] = {
49-
val cfg = config.getConfig(s"com.sandinh.couchbase.buckets.$bucket")
50-
val name = Try { cfg.getString("name") } getOrElse bucket
51-
val pass = cfg.getString("password")
59+
val cfg = config.getConfig("com.sandinh.couchbase")
60+
val name = s"buckets.$bucket.name" match {
61+
case p if cfg.hasPath(p) => cfg.getString(p)
62+
case _ => bucket
63+
}
5264
val stringTranscoder =
5365
if (legacyEncodeString) CompatStringTranscoderLegacy
5466
else CompatStringTranscoder
5567
val trans = transcoders :+ JsTranscoder :+ stringTranscoder
56-
asJava.openBucket(name, pass, trans.asJava).scMap(_.asScala).toFuture
68+
val bucketObs = asJava.authenticator() match {
69+
case _: PasswordAuthenticator | _: CertAuthenticator =>
70+
asJava.openBucket(name, trans.asJava)
71+
case _ =>
72+
val pass = cfg.getString(s"buckets.$bucket.password")
73+
asJava.openBucket(name, pass, trans.asJava)
74+
}
75+
bucketObs.scMap(_.asScala).toFuture
5776
}
5877

5978
/** @note You should never perform long-running blocking operations inside of an asynchronous stream (e.g. inside of maps or flatMaps)

0 commit comments

Comments
 (0)