Skip to content

Commit 6b6b6f9

Browse files
flowchartsmanmichaeljmarshall
authored andcommitted
[fix][fn] enable Go function token auth and TLS (apache#20468)
Note: there were conflicts with the go.sum file. In order to get things working, I cleared out the conflicts and then ran `go mod tidy` in the `pulsar-funcion-go/examples` directory. Co-authored-by: Andy Walker <andy@andy.dev> (cherry picked from commit 8b3c085)
1 parent 0c40723 commit 6b6b6f9

8 files changed

Lines changed: 109 additions & 61 deletions

File tree

pulsar-function-go/conf/conf.go

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -27,7 +27,6 @@ import (
2727
"time"
2828

2929
log "github.com/apache/pulsar/pulsar-function-go/logutil"
30-
"gopkg.in/yaml.v2"
3130
)
3231

3332
const ConfigPath = "conf/conf.yaml"
@@ -51,6 +50,12 @@ type Conf struct {
5150
Runtime int32 `json:"runtime" yaml:"runtime"`
5251
AutoACK bool `json:"autoAck" yaml:"autoAck"`
5352
Parallelism int32 `json:"parallelism" yaml:"parallelism"`
53+
// Authentication
54+
ClientAuthenticationPlugin string `json:"clientAuthenticationPlugin" yaml:"clientAuthenticationPlugin"`
55+
ClientAuthenticationParameters string `json:"clientAuthenticationParameters" yaml:"clientAuthenticationParameters"`
56+
TLSTrustCertsFilePath string `json:"tlsTrustCertsFilePath" yaml:"tlsTrustCertsFilePath"`
57+
TLSAllowInsecureConnection bool `json:"tlsAllowInsecureConnection" yaml:"tlsAllowInsecureConnection"`
58+
TLSHostnameVerificationEnable bool `json:"tlsHostnameVerificationEnable" yaml:"tlsHostnameVerificationEnable"`
5459
//source config
5560
SubscriptionType int32 `json:"subscriptionType" yaml:"subscriptionType"`
5661
TimeoutMs uint64 `json:"timeoutMs" yaml:"timeoutMs"`

pulsar-function-go/examples/go.mod

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -5,8 +5,6 @@ go 1.13
55
require (
66
github.com/apache/pulsar-client-go v0.7.0
77
github.com/apache/pulsar/pulsar-function-go v0.0.0
8-
github.com/datadog/zstd v1.4.6-0.20200617134701-89f69fb7df32 // indirect
9-
github.com/yahoo/athenz v1.8.55 // indirect
108
)
119

1210
replace github.com/apache/pulsar/pulsar-function-go => ../

pulsar-function-go/examples/go.sum

Lines changed: 19 additions & 53 deletions
Large diffs are not rendered by default.

pulsar-function-go/pf/instance.go

Lines changed: 34 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -21,8 +21,10 @@ package pf
2121

2222
import (
2323
"context"
24+
"fmt"
2425
"math"
2526
"strconv"
27+
"strings"
2628
"time"
2729

2830
"github.com/golang/protobuf/ptypes/empty"
@@ -193,11 +195,40 @@ CLOSE:
193195
return nil
194196
}
195197

198+
const (
199+
authPluginToken = "org.apache.pulsar.client.impl.auth.AuthenticationToken"
200+
authPluginNone = ""
201+
)
202+
196203
func (gi *goInstance) setupClient() error {
197-
client, err := pulsar.NewClient(pulsar.ClientOptions{
204+
ic := gi.context.instanceConf
205+
206+
clientOpts := pulsar.ClientOptions{
207+
URL: ic.pulsarServiceURL,
208+
TLSTrustCertsFilePath: ic.tlsTrustCertsPath,
209+
TLSAllowInsecureConnection: ic.tlsAllowInsecure,
210+
TLSValidateHostname: ic.tlsHostnameVerification,
211+
}
212+
213+
switch ic.authPlugin {
214+
case authPluginToken:
215+
switch {
216+
case strings.HasPrefix(ic.authParams, "file://"):
217+
clientOpts.Authentication = pulsar.NewAuthenticationTokenFromFile(ic.authParams[7:])
218+
case strings.HasPrefix(ic.authParams, "token:"):
219+
clientOpts.Authentication = pulsar.NewAuthenticationToken(ic.authParams[6:])
220+
case ic.authParams == "":
221+
return fmt.Errorf("auth plugin %s given, but authParams is empty", authPluginToken)
222+
default:
223+
return fmt.Errorf(`unknown token format - expecting "file://" or "token:" prefix`)
224+
}
225+
case authPluginNone:
226+
clientOpts.Authentication, _ = pulsar.NewAuthentication("", "") // ret: auth.NewAuthDisabled()
227+
default:
228+
return fmt.Errorf("unknown auth provider: %s", ic.authPlugin)
229+
}
198230

199-
URL: gi.context.instanceConf.pulsarServiceURL,
200-
})
231+
client, err := pulsar.NewClient(clientOpts)
201232
if err != nil {
202233
log.Errorf("create client error:%v", err)
203234
gi.stats.incrTotalSysExceptions(err)

pulsar-function-go/pf/instanceConf.go

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -41,6 +41,11 @@ type instanceConf struct {
4141
killAfterIdle time.Duration
4242
expectedHealthCheckInterval int32
4343
metricsPort int
44+
authPlugin string
45+
authParams string
46+
tlsTrustCertsPath string
47+
tlsAllowInsecure bool
48+
tlsHostnameVerification bool
4449
}
4550

4651
func newInstanceConf() *instanceConf {
@@ -102,6 +107,11 @@ func newInstanceConf() *instanceConf {
102107
},
103108
UserConfig: cfg.UserConfig,
104109
},
110+
authPlugin: cfg.ClientAuthenticationPlugin,
111+
authParams: cfg.ClientAuthenticationParameters,
112+
tlsTrustCertsPath: cfg.TLSTrustCertsFilePath,
113+
tlsAllowInsecure: cfg.TLSAllowInsecureConnection,
114+
tlsHostnameVerification: cfg.TLSHostnameVerificationEnable,
105115
}
106116
return instanceConf
107117
}

pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/instance/go/GoInstanceConfig.java

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -43,6 +43,13 @@ public class GoInstanceConfig {
4343
private int processingGuarantees;
4444
private String secretsMap = "";
4545
private String userConfig = "";
46+
47+
private String clientAuthenticationPlugin = "";
48+
private String clientAuthenticationParameters = "";
49+
private String tlsTrustCertsFilePath = "";
50+
private boolean tlsHostnameVerificationEnable = false;
51+
private boolean tlsAllowInsecureConnection = false;
52+
4653
private int runtime;
4754
private boolean autoAck;
4855
private int parallelism;

pulsar-functions/runtime/src/main/java/org/apache/pulsar/functions/runtime/RuntimeUtils.java

Lines changed: 20 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -133,6 +133,7 @@ public static List<String> getArgsBeforeCmd(InstanceConfig instanceConfig, Strin
133133
*/
134134

135135
public static List<String> getGoInstanceCmd(InstanceConfig instanceConfig,
136+
AuthenticationConfig authConfig,
136137
String originalCodeFileName,
137138
String pulsarServiceUrl,
138139
boolean k8sRuntime) throws IOException {
@@ -190,6 +191,23 @@ public static List<String> getGoInstanceCmd(InstanceConfig instanceConfig,
190191
goInstanceConfig.setParallelism(instanceConfig.getFunctionDetails().getParallelism());
191192
}
192193

194+
if (authConfig != null) {
195+
if (isNotBlank(authConfig.getClientAuthenticationPlugin())
196+
&& isNotBlank(authConfig.getClientAuthenticationParameters())) {
197+
goInstanceConfig.setClientAuthenticationPlugin(authConfig.getClientAuthenticationPlugin());
198+
goInstanceConfig.setClientAuthenticationParameters(authConfig.getClientAuthenticationParameters());
199+
}
200+
goInstanceConfig.setTlsAllowInsecureConnection(
201+
authConfig.isTlsAllowInsecureConnection());
202+
goInstanceConfig.setTlsHostnameVerificationEnable(
203+
authConfig.isTlsHostnameVerificationEnable());
204+
if (isNotBlank(authConfig.getTlsTrustCertsFilePath())){
205+
goInstanceConfig.setTlsTrustCertsFilePath(
206+
authConfig.getTlsTrustCertsFilePath());
207+
}
208+
209+
}
210+
193211
if (instanceConfig.getMaxBufferedTuples() != 0) {
194212
goInstanceConfig.setMaxBufTuples(instanceConfig.getMaxBufferedTuples());
195213
}
@@ -286,7 +304,8 @@ public static List<String> getCmd(InstanceConfig instanceConfig,
286304
final List<String> args = new LinkedList<>();
287305

288306
if (instanceConfig.getFunctionDetails().getRuntime() == Function.FunctionDetails.Runtime.GO) {
289-
return getGoInstanceCmd(instanceConfig, originalCodeFileName, pulsarServiceUrl, k8sRuntime);
307+
return getGoInstanceCmd(instanceConfig, authConfig,
308+
originalCodeFileName, pulsarServiceUrl, k8sRuntime);
290309
}
291310

292311
if (instanceConfig.getFunctionDetails().getRuntime() == Function.FunctionDetails.Runtime.JAVA) {

pulsar-functions/runtime/src/test/java/org/apache/pulsar/functions/runtime/RuntimeUtilsTest.java

Lines changed: 13 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -73,6 +73,13 @@ public void getGoInstanceCmd(boolean k8sRuntime) throws IOException {
7373
instanceConfig.setPort(1337);
7474
instanceConfig.setMetricsPort(60000);
7575

76+
AuthenticationConfig authConfig = AuthenticationConfig.builder()
77+
.clientAuthenticationPlugin("org.apache.pulsar.client.impl.auth.AuthenticationToken")
78+
.clientAuthenticationParameters("file:///secret/token.jwt")
79+
.tlsTrustCertsFilePath("/secret/ca.cert.pem")
80+
.tlsHostnameVerificationEnable(true)
81+
.tlsAllowInsecureConnection(false)
82+
.build();
7683

7784
JSONObject userConfig = new JSONObject();
7885
userConfig.put("word-of-the-day", "der Weltschmerz");
@@ -116,7 +123,7 @@ public void getGoInstanceCmd(boolean k8sRuntime) throws IOException {
116123

117124
instanceConfig.setFunctionDetails(functionDetails);
118125

119-
List<String> commands = RuntimeUtils.getGoInstanceCmd(instanceConfig, "config", "pulsar://localhost:6650", k8sRuntime);
126+
List<String> commands = RuntimeUtils.getGoInstanceCmd(instanceConfig, authConfig,"config", "pulsar://localhost:6650", k8sRuntime);
120127
if (k8sRuntime) {
121128
goInstanceConfig = new ObjectMapper().readValue(commands.get(2).replaceAll("^\'|\'$", ""), HashMap.class);
122129
} else {
@@ -160,6 +167,11 @@ public void getGoInstanceCmd(boolean k8sRuntime) throws IOException {
160167
Assert.assertEquals(goInstanceConfig.get("deadLetterTopic"), "go-func-deadletter");
161168
Assert.assertEquals(goInstanceConfig.get("userConfig"), userConfig.toString());
162169
Assert.assertEquals(goInstanceConfig.get("metricsPort"), 60000);
170+
Assert.assertEquals(goInstanceConfig.get("clientAuthenticationPlugin"), "org.apache.pulsar.client.impl.auth.AuthenticationToken");
171+
Assert.assertEquals(goInstanceConfig.get("clientAuthenticationParameters"), "file:///secret/token.jwt");
172+
Assert.assertEquals(goInstanceConfig.get("tlsTrustCertsFilePath"), "/secret/ca.cert.pem");
173+
Assert.assertEquals(goInstanceConfig.get("tlsHostnameVerificationEnable"), true);
174+
Assert.assertEquals(goInstanceConfig.get("tlsAllowInsecureConnection"), false);
163175
}
164176

165177
@DataProvider(name = "k8sRuntime")

0 commit comments

Comments
 (0)