From 90b922a003bf0b74930923deffae55f3773680cb Mon Sep 17 00:00:00 2001 From: Zhu Zhanyan Date: Thu, 6 Aug 2020 18:50:04 +0800 Subject: [PATCH 01/18] Fix race condition in GrpcMessageInterceptor to revert a empty message if message cannot be recorded. --- .../common/logging/interceptors/GrpcMessageInterceptor.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/common/src/main/java/feast/common/logging/interceptors/GrpcMessageInterceptor.java b/common/src/main/java/feast/common/logging/interceptors/GrpcMessageInterceptor.java index da43a372d73..6f3bfbb45b5 100644 --- a/common/src/main/java/feast/common/logging/interceptors/GrpcMessageInterceptor.java +++ b/common/src/main/java/feast/common/logging/interceptors/GrpcMessageInterceptor.java @@ -61,7 +61,8 @@ public GrpcMessageInterceptor(@Nullable SecurityProperties securityProperties) { public Listener interceptCall( ServerCall call, Metadata headers, ServerCallHandler next) { MessageAuditLogEntry.Builder entryBuilder = MessageAuditLogEntry.newBuilder(); - // default response message to empty proto in log entry. + // default request/response message to empty proto in log entry. + entryBuilder.setRequest(Empty.newBuilder().build()); entryBuilder.setResponse(Empty.newBuilder().build()); // Unpack service & method name from call From 88f2210d9cbe1f201f81fb446f06877d6af85c56 Mon Sep 17 00:00:00 2001 From: Zhu Zhanyan Date: Fri, 7 Aug 2020 10:27:17 +0800 Subject: [PATCH 02/18] Fix GrpcMessageInterceptor race condition by allow audit log message to be called from multiple async calls --- .../logging/entry/MessageAuditLogEntry.java | 22 ++++++++++++++++++- .../interceptors/GrpcMessageInterceptor.java | 19 ++++++++-------- core/src/main/resources/application.yml | 2 +- 3 files changed, 31 insertions(+), 12 deletions(-) diff --git a/common/src/main/java/feast/common/logging/entry/MessageAuditLogEntry.java b/common/src/main/java/feast/common/logging/entry/MessageAuditLogEntry.java index 745cc1283ae..f158dae4221 100644 --- a/common/src/main/java/feast/common/logging/entry/MessageAuditLogEntry.java +++ b/common/src/main/java/feast/common/logging/entry/MessageAuditLogEntry.java @@ -66,24 +66,44 @@ public abstract class MessageAuditLogEntry extends AuditLogEntry { public abstract static class Builder { public abstract Builder setId(UUID id); + public abstract UUID getId(); + public abstract Builder setComponent(String component); - public abstract Builder setVersion(String component); + public abstract String getComponent(); + + public abstract Builder setVersion(String version); + + public abstract String getVersion(); public abstract Builder setKind(AuditLogEntryKind kind); + public abstract AuditLogEntryKind getKind(); + public abstract Builder setService(String name); + public abstract String getService(); + public abstract Builder setMethod(String name); + public abstract String getMethod(); + public abstract Builder setRequest(Message request); + public abstract Message getRequest(); + public abstract Builder setResponse(Message response); + public abstract Message getResponse(); + public abstract Builder setIdentity(String identity); + public abstract String getIdentity(); + public abstract Builder setStatusCode(Code statusCode); + public abstract Code getStatusCode(); + public abstract MessageAuditLogEntry build(); } diff --git a/common/src/main/java/feast/common/logging/interceptors/GrpcMessageInterceptor.java b/common/src/main/java/feast/common/logging/interceptors/GrpcMessageInterceptor.java index 6f3bfbb45b5..2ea909f00c4 100644 --- a/common/src/main/java/feast/common/logging/interceptors/GrpcMessageInterceptor.java +++ b/common/src/main/java/feast/common/logging/interceptors/GrpcMessageInterceptor.java @@ -16,7 +16,6 @@ */ package feast.common.logging.interceptors; -import com.google.protobuf.Empty; import com.google.protobuf.Message; import feast.common.auth.config.SecurityProperties; import feast.common.auth.config.SecurityProperties.AuthenticationProperties; @@ -61,10 +60,6 @@ public GrpcMessageInterceptor(@Nullable SecurityProperties securityProperties) { public Listener interceptCall( ServerCall call, Metadata headers, ServerCallHandler next) { MessageAuditLogEntry.Builder entryBuilder = MessageAuditLogEntry.newBuilder(); - // default request/response message to empty proto in log entry. - entryBuilder.setRequest(Empty.newBuilder().build()); - entryBuilder.setResponse(Empty.newBuilder().build()); - // Unpack service & method name from call // full method name is in format ./ String fullMethodName = call.getMethodDescriptor().getFullMethodName(); @@ -78,22 +73,25 @@ public Listener interceptCall( entryBuilder.setIdentity(identity); // Register forwarding call to intercept outgoing response and log to audit log + // As sendMessage(), close() and onMessage() may not be called in order (async) + // add a try log message call to each so that message logging would not + // depend on the calls to made in a specific order. call = new SimpleForwardingServerCall(call) { @Override public void sendMessage(RespT message) { - // 2. Track the response & Log entry to audit logger + // Track the response & Log entry to audit logger super.sendMessage(message); entryBuilder.setResponse((Message) message); + tryLogMessage(entryBuilder); } @Override public void close(Status status, Metadata trailers) { super.close(status, trailers); - // 3. Log the message log entry to the audit log - Level logLevel = (status.isOk()) ? Level.INFO : Level.ERROR; + // Log the message log entry to the audit log entryBuilder.setStatusCode(status.getCode()); - AuditLogger.logMessage(logLevel, entryBuilder); + tryLogMessage(entryBuilder); } }; @@ -103,8 +101,9 @@ public void close(Status status, Metadata trailers) { // Register listener to intercept incoming request messages and log to audit log public void onMessage(ReqT message) { super.onMessage(message); - // 1. Track the request. + // Track the request. entryBuilder.setRequest((Message) message); + tryLogMessage(entryBuilder); } }; } diff --git a/core/src/main/resources/application.yml b/core/src/main/resources/application.yml index 15d3b9bd3a5..7dff5e90eb7 100644 --- a/core/src/main/resources/application.yml +++ b/core/src/main/resources/application.yml @@ -57,7 +57,7 @@ feast: # Whether audit logging is enabled. enabled: true # Whether to enable message level (ie request/response) audit logging - messageLoggingEnabled: false + messageLoggingEnabled: true grpc: server: From e8861a6ccb64b319a316758327078b40817d013c Mon Sep 17 00:00:00 2001 From: Zhu Zhanyan Date: Fri, 7 Aug 2020 15:32:46 +0800 Subject: [PATCH 03/18] Revert option change --- core/src/main/resources/application.yml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/core/src/main/resources/application.yml b/core/src/main/resources/application.yml index 7dff5e90eb7..15d3b9bd3a5 100644 --- a/core/src/main/resources/application.yml +++ b/core/src/main/resources/application.yml @@ -57,7 +57,7 @@ feast: # Whether audit logging is enabled. enabled: true # Whether to enable message level (ie request/response) audit logging - messageLoggingEnabled: true + messageLoggingEnabled: false grpc: server: From 282540796998bac87aeb6c25669ae8c1aa65ff3a Mon Sep 17 00:00:00 2001 From: Zhu Zhanyan Date: Wed, 19 Aug 2020 17:24:11 +0800 Subject: [PATCH 04/18] Add CoreLoggingIT integration test to test message audit logging. --- .../interceptors/GrpcMessageInterceptor.java | 25 ++- .../feast/core/logging/CoreLoggingIT.java | 171 ++++++++++++++++++ .../feast/core/logging/TestLogAppender.java | 63 +++++++ core/src/test/resources/log4j2-test.xml | 53 ++++++ 4 files changed, 303 insertions(+), 9 deletions(-) create mode 100644 core/src/test/java/feast/core/logging/CoreLoggingIT.java create mode 100644 core/src/test/java/feast/core/logging/TestLogAppender.java create mode 100644 core/src/test/resources/log4j2-test.xml diff --git a/common/src/main/java/feast/common/logging/interceptors/GrpcMessageInterceptor.java b/common/src/main/java/feast/common/logging/interceptors/GrpcMessageInterceptor.java index 2ea909f00c4..83a6d00bf6f 100644 --- a/common/src/main/java/feast/common/logging/interceptors/GrpcMessageInterceptor.java +++ b/common/src/main/java/feast/common/logging/interceptors/GrpcMessageInterceptor.java @@ -16,6 +16,7 @@ */ package feast.common.logging.interceptors; +import com.google.protobuf.Empty; import com.google.protobuf.Message; import feast.common.auth.config.SecurityProperties; import feast.common.auth.config.SecurityProperties.AuthenticationProperties; @@ -59,39 +60,45 @@ public GrpcMessageInterceptor(@Nullable SecurityProperties securityProperties) { @Override public Listener interceptCall( ServerCall call, Metadata headers, ServerCallHandler next) { + // System.out.println(">>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>HERE>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>"); MessageAuditLogEntry.Builder entryBuilder = MessageAuditLogEntry.newBuilder(); + // default response message to empty proto in log entry. + entryBuilder.setResponse(Empty.newBuilder().build()); + // System.out.println(">>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>Response>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>"); + // Unpack service & method name from call // full method name is in format ./ String fullMethodName = call.getMethodDescriptor().getFullMethodName(); entryBuilder.setService( fullMethodName.substring(fullMethodName.lastIndexOf(".") + 1, fullMethodName.indexOf("/"))); entryBuilder.setMethod(fullMethodName.substring(fullMethodName.indexOf("/") + 1)); + // System.out.println(">>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>Method>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>"); // Attempt Extract current authenticated identity. Authentication authentication = SecurityContextHolder.getContext().getAuthentication(); String identity = (authentication != null) ? getIdentity(authentication) : ""; entryBuilder.setIdentity(identity); + // System.out.println(">>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>AUTH>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>"); // Register forwarding call to intercept outgoing response and log to audit log - // As sendMessage(), close() and onMessage() may not be called in order (async) - // add a try log message call to each so that message logging would not - // depend on the calls to made in a specific order. call = new SimpleForwardingServerCall(call) { @Override public void sendMessage(RespT message) { - // Track the response & Log entry to audit logger + // 2. Track the response & Log entry to audit logger super.sendMessage(message); entryBuilder.setResponse((Message) message); - tryLogMessage(entryBuilder); + // System.out.println(">>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>OUTGOING>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>"); } @Override public void close(Status status, Metadata trailers) { super.close(status, trailers); - // Log the message log entry to the audit log + // 3. Log the message log entry to the audit log + Level logLevel = (status.isOk()) ? Level.INFO : Level.ERROR; entryBuilder.setStatusCode(status.getCode()); - tryLogMessage(entryBuilder); + AuditLogger.logMessage(logLevel, entryBuilder); + // System.out.println(">>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>END>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>"); } }; @@ -101,9 +108,9 @@ public void close(Status status, Metadata trailers) { // Register listener to intercept incoming request messages and log to audit log public void onMessage(ReqT message) { super.onMessage(message); - // Track the request. + // 1. Track the request. entryBuilder.setRequest((Message) message); - tryLogMessage(entryBuilder); + // System.out.println(">>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>INCOMING>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>"); } }; } diff --git a/core/src/test/java/feast/core/logging/CoreLoggingIT.java b/core/src/test/java/feast/core/logging/CoreLoggingIT.java new file mode 100644 index 00000000000..3753cedc41a --- /dev/null +++ b/core/src/test/java/feast/core/logging/CoreLoggingIT.java @@ -0,0 +1,171 @@ +/* + * SPDX-License-Identifier: Apache-2.0 + * Copyright 2018-2020 The Feast Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package feast.core.logging; + +import static org.hamcrest.CoreMatchers.*; +import static org.hamcrest.MatcherAssert.assertThat; +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertTrue; + +import com.google.gson.JsonObject; +import com.google.gson.JsonParser; +import com.google.protobuf.InvalidProtocolBufferException; +import com.google.protobuf.util.JsonFormat; +import feast.common.logging.entry.AuditLogEntryKind; +import feast.core.it.BaseIT; +import feast.core.it.DataGenerator; +import feast.proto.core.CoreServiceGrpc; +import feast.proto.core.CoreServiceGrpc.CoreServiceBlockingStub; +import feast.proto.core.CoreServiceProto.GetFeastCoreVersionRequest; +import feast.proto.core.CoreServiceProto.ListFeatureSetsRequest; +import feast.proto.core.CoreServiceProto.UpdateStoreRequest; +import feast.proto.core.CoreServiceProto.UpdateStoreResponse; +import io.grpc.Channel; +import io.grpc.ManagedChannelBuilder; +import io.grpc.StatusRuntimeException; +import io.grpc.Status.Code; + +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.ExecutionException; +import java.util.stream.Collectors; +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.core.LoggerContext; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.Test; +import org.slf4j.event.Level; +import org.springframework.beans.factory.annotation.Value; +import org.springframework.boot.test.context.SpringBootTest; + +@SpringBootTest( + properties = { + "feast.logging.audit.enabled=true", + "feast.logging.audit.messageLoggingEnabled=true", + }) +public class CoreLoggingIT extends BaseIT { + private static TestLogAppender testAuditLogAppender; + private static CoreServiceBlockingStub coreService; + + @BeforeAll + public static void globalSetUp(@Value("${grpc.server.port}") int coreGrpcPort) + throws InterruptedException, ExecutionException { + LoggerContext logContext = (LoggerContext) LogManager.getContext(false); + testAuditLogAppender = logContext.getConfiguration().getAppender("TestAuditLogAppender"); + + // Connect to core. + Channel channel = + ManagedChannelBuilder.forAddress("localhost", coreGrpcPort).usePlaintext().build(); + coreService = CoreServiceGrpc.newBlockingStub(channel); + + // Preflight a request to core service stub to verify connection + coreService.getFeastCoreVersion(GetFeastCoreVersionRequest.getDefaultInstance()); + } + + /** + * Test Message Audit Logging. Produce artifical load on Feast core and checking the produced + * message logs. + */ + @Test + public void shouldProduceMessageAuditLogs() throws InterruptedException { + testAuditLogAppender.getLogs().clear(); + // Generate artifical load on feast core. + UpdateStoreRequest request = + UpdateStoreRequest.newBuilder().setStore(DataGenerator.getDefaultStore()).build(); + UpdateStoreResponse receivedResponse = null; + final int LOAD_SIZE = 200; + final int ASYNC_BURST_SIZE = 5; + for (int i = 0; i < LOAD_SIZE; i++) { + receivedResponse = coreService.updateStore(request); + } + final UpdateStoreResponse response = receivedResponse; + + // Wait required to ensure audit logs are flushed into test audit log appender + Thread.sleep(500); + // Check message audit logs are produced for each audit log. + JsonFormat.Parser protoJSONParser = JsonFormat.parser(); + List logJsonObjects = parseMessageJsonLogObjects(testAuditLogAppender.getLogs()); + + logJsonObjects.forEach( + logObj -> { + // Extract recorded request/response from message logs + try { + String requestJson = logObj.getAsJsonObject("request").toString(); + UpdateStoreRequest.Builder gotRequest = UpdateStoreRequest.newBuilder(); + protoJSONParser.merge(requestJson, gotRequest); + + String responseJson = logObj.getAsJsonObject("response").toString(); + UpdateStoreResponse.Builder gotResponse = UpdateStoreResponse.newBuilder(); + protoJSONParser.merge(responseJson, gotResponse); + + // Check that request/response are returned correctlyf + assertThat(gotRequest.build(), equalTo(request)); + assertThat(gotResponse.build(), equalTo(response)); + } catch (InvalidProtocolBufferException e) { + throw new RuntimeException(e); + } + }); + } + + /** Check that message audit logs are produced when server encounters an error */ + @Test + public void shouldProduceMessageAuditLogsOnError() throws InterruptedException { + testAuditLogAppender.getLogs().clear(); + // send a bad request which should cause Core to error + ListFeatureSetsRequest request = + ListFeatureSetsRequest.newBuilder() + .setFilter( + ListFeatureSetsRequest.Filter.newBuilder() + .setProject("*") + .setFeatureSetName("nop") + .build()) + .build(); + + boolean hasExpectedException = false; + Code statusCode = null; + try { + coreService.listFeatureSets(request); + } catch (StatusRuntimeException e) { + hasExpectedException = true; + statusCode = e.getStatus().getCode(); + } + assertTrue(hasExpectedException); + + // Wait required to ensure audit logs are flushed into test audit log appender + Thread.sleep(500); + List logJsonObjects = parseMessageJsonLogObjects(testAuditLogAppender.getLogs()); + assertEquals(logJsonObjects.size(), 1); + JsonObject logJsonObject = logJsonObjects.get(0); + assertEquals(logJsonObject.get("statusCode").getAsString(), statusCode.toString()); + } + + /** Filter and Parse out Message Audit Logs from the given logsStrings */ + private List parseMessageJsonLogObjects(List logsStrings) { + JsonParser jsonParser = new JsonParser(); + // copy to prevent concurrent modification. + List logsStringsCopy = new ArrayList(logsStrings); + return logsStringsCopy.stream() + .map(logJSON -> jsonParser.parse(logJSON).getAsJsonObject()) + // Filter to only include message audit logs + .filter( + logObj -> + logObj + .getAsJsonPrimitive("kind") + .getAsString() + .equals(AuditLogEntryKind.MESSAGE.toString())) + .collect(Collectors.toList()); + } +} diff --git a/core/src/test/java/feast/core/logging/TestLogAppender.java b/core/src/test/java/feast/core/logging/TestLogAppender.java new file mode 100644 index 00000000000..21c5bbd69d0 --- /dev/null +++ b/core/src/test/java/feast/core/logging/TestLogAppender.java @@ -0,0 +1,63 @@ +/* + * SPDX-License-Identifier: Apache-2.0 + * Copyright 2018-2020 The Feast Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package feast.core.logging; + +import java.io.Serializable; +import java.util.ArrayList; +import java.util.List; +import lombok.Getter; +import org.apache.logging.log4j.core.Filter; +import org.apache.logging.log4j.core.Layout; +import org.apache.logging.log4j.core.LogEvent; +import org.apache.logging.log4j.core.appender.AbstractAppender; +import org.apache.logging.log4j.core.config.Property; +import org.apache.logging.log4j.core.config.plugins.Plugin; +import org.apache.logging.log4j.core.config.plugins.PluginAttribute; +import org.apache.logging.log4j.core.config.plugins.PluginElement; +import org.apache.logging.log4j.core.config.plugins.PluginFactory; +import org.apache.logging.log4j.core.layout.PatternLayout; + +/** Test Log Appender used for collecting logs for testing logging. */ +@Plugin(name = "TestLogAppender", category = "Core", elementType = "appender", printObject = true) +@Getter +public class TestLogAppender extends AbstractAppender { + private List logs; + + protected TestLogAppender(String name, Filter filter, Layout layout) { + super(name, filter, layout, false, new Property[] {}); + logs = new ArrayList<>(); + } + + @Override + public void append(LogEvent event) { + getLogs().add(event.getMessage().toString()); + } + + @PluginFactory + public static TestLogAppender createAppender( + @PluginAttribute("name") String name, + @PluginElement("Layout") Layout layout, + @PluginElement("Filter") final Filter filter) { + if (name == null) { + return null; + } + if (layout == null) { + layout = PatternLayout.createDefaultLayout(); + } + return new TestLogAppender(name, filter, layout); + } +} diff --git a/core/src/test/resources/log4j2-test.xml b/core/src/test/resources/log4j2-test.xml new file mode 100644 index 00000000000..90945373bfe --- /dev/null +++ b/core/src/test/resources/log4j2-test.xml @@ -0,0 +1,53 @@ + + + + + + + %d{yyyy-MM-dd HH:mm:ss.SSS} %5p ${hostName} --- [%15.15t] %-40.40c{1.} : %m%n%ex + + + {"time":"%d{yyyy-MM-dd'T'HH:mm:ssXXX}","hostname":"${hostName}","severity":"%p","message":%m}%n%ex + + + + + + + + + + + + + + + + + + + + + + + + + + + + From db3e5dcca6990be31082dee6da7b13b7a2b49d97 Mon Sep 17 00:00:00 2001 From: Zhu Zhanyan Date: Wed, 19 Aug 2020 17:27:26 +0800 Subject: [PATCH 05/18] Fix GrpcMessageInterceptor race condition by moving allowing request to be unset. --- .../logging/interceptors/GrpcMessageInterceptor.java | 12 ++++-------- 1 file changed, 4 insertions(+), 8 deletions(-) diff --git a/common/src/main/java/feast/common/logging/interceptors/GrpcMessageInterceptor.java b/common/src/main/java/feast/common/logging/interceptors/GrpcMessageInterceptor.java index 83a6d00bf6f..f38c6fdfe89 100644 --- a/common/src/main/java/feast/common/logging/interceptors/GrpcMessageInterceptor.java +++ b/common/src/main/java/feast/common/logging/interceptors/GrpcMessageInterceptor.java @@ -60,11 +60,12 @@ public GrpcMessageInterceptor(@Nullable SecurityProperties securityProperties) { @Override public Listener interceptCall( ServerCall call, Metadata headers, ServerCallHandler next) { - // System.out.println(">>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>HERE>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>"); MessageAuditLogEntry.Builder entryBuilder = MessageAuditLogEntry.newBuilder(); - // default response message to empty proto in log entry. + // default response/request message to empty proto in log entry. + // request could be empty when the client closes the connection before sending a request message. + // response could be unset when the service encounters an error when processsing the service call. + entryBuilder.setRequest(Empty.newBuilder().build()); entryBuilder.setResponse(Empty.newBuilder().build()); - // System.out.println(">>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>Response>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>"); // Unpack service & method name from call // full method name is in format ./ @@ -72,13 +73,11 @@ public Listener interceptCall( entryBuilder.setService( fullMethodName.substring(fullMethodName.lastIndexOf(".") + 1, fullMethodName.indexOf("/"))); entryBuilder.setMethod(fullMethodName.substring(fullMethodName.indexOf("/") + 1)); - // System.out.println(">>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>Method>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>"); // Attempt Extract current authenticated identity. Authentication authentication = SecurityContextHolder.getContext().getAuthentication(); String identity = (authentication != null) ? getIdentity(authentication) : ""; entryBuilder.setIdentity(identity); - // System.out.println(">>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>AUTH>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>"); // Register forwarding call to intercept outgoing response and log to audit log call = @@ -88,7 +87,6 @@ public void sendMessage(RespT message) { // 2. Track the response & Log entry to audit logger super.sendMessage(message); entryBuilder.setResponse((Message) message); - // System.out.println(">>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>OUTGOING>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>"); } @Override @@ -98,7 +96,6 @@ public void close(Status status, Metadata trailers) { Level logLevel = (status.isOk()) ? Level.INFO : Level.ERROR; entryBuilder.setStatusCode(status.getCode()); AuditLogger.logMessage(logLevel, entryBuilder); - // System.out.println(">>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>END>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>"); } }; @@ -110,7 +107,6 @@ public void onMessage(ReqT message) { super.onMessage(message); // 1. Track the request. entryBuilder.setRequest((Message) message); - // System.out.println(">>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>INCOMING>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>>"); } }; } From 42fe1dc0e3d64c1f49f1503e9d24bd71a42f1c3e Mon Sep 17 00:00:00 2001 From: Zhu Zhanyan Date: Wed, 19 Aug 2020 17:33:32 +0800 Subject: [PATCH 06/18] Fix lint --- .../common/logging/interceptors/GrpcMessageInterceptor.java | 6 ++++-- core/src/test/java/feast/core/logging/CoreLoggingIT.java | 4 +--- 2 files changed, 5 insertions(+), 5 deletions(-) diff --git a/common/src/main/java/feast/common/logging/interceptors/GrpcMessageInterceptor.java b/common/src/main/java/feast/common/logging/interceptors/GrpcMessageInterceptor.java index f38c6fdfe89..abfca86ecde 100644 --- a/common/src/main/java/feast/common/logging/interceptors/GrpcMessageInterceptor.java +++ b/common/src/main/java/feast/common/logging/interceptors/GrpcMessageInterceptor.java @@ -62,8 +62,10 @@ public Listener interceptCall( ServerCall call, Metadata headers, ServerCallHandler next) { MessageAuditLogEntry.Builder entryBuilder = MessageAuditLogEntry.newBuilder(); // default response/request message to empty proto in log entry. - // request could be empty when the client closes the connection before sending a request message. - // response could be unset when the service encounters an error when processsing the service call. + // request could be empty when the client closes the connection before sending a request + // message. + // response could be unset when the service encounters an error when processsing the service + // call. entryBuilder.setRequest(Empty.newBuilder().build()); entryBuilder.setResponse(Empty.newBuilder().build()); diff --git a/core/src/test/java/feast/core/logging/CoreLoggingIT.java b/core/src/test/java/feast/core/logging/CoreLoggingIT.java index 3753cedc41a..841eada64a2 100644 --- a/core/src/test/java/feast/core/logging/CoreLoggingIT.java +++ b/core/src/test/java/feast/core/logging/CoreLoggingIT.java @@ -36,9 +36,8 @@ import feast.proto.core.CoreServiceProto.UpdateStoreResponse; import io.grpc.Channel; import io.grpc.ManagedChannelBuilder; -import io.grpc.StatusRuntimeException; import io.grpc.Status.Code; - +import io.grpc.StatusRuntimeException; import java.util.ArrayList; import java.util.List; import java.util.concurrent.ExecutionException; @@ -47,7 +46,6 @@ import org.apache.logging.log4j.core.LoggerContext; import org.junit.jupiter.api.BeforeAll; import org.junit.jupiter.api.Test; -import org.slf4j.event.Level; import org.springframework.beans.factory.annotation.Value; import org.springframework.boot.test.context.SpringBootTest; From 469b647d5fc3664f2ca1d6d118c3f1aefe66ce05 Mon Sep 17 00:00:00 2001 From: Zhu Zhanyan Date: Wed, 19 Aug 2020 18:22:56 +0800 Subject: [PATCH 07/18] Increase wait for logs in CoreLoggingIT --- core/src/test/java/feast/core/logging/CoreLoggingIT.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/core/src/test/java/feast/core/logging/CoreLoggingIT.java b/core/src/test/java/feast/core/logging/CoreLoggingIT.java index 841eada64a2..161c0e257f9 100644 --- a/core/src/test/java/feast/core/logging/CoreLoggingIT.java +++ b/core/src/test/java/feast/core/logging/CoreLoggingIT.java @@ -92,7 +92,7 @@ public void shouldProduceMessageAuditLogs() throws InterruptedException { final UpdateStoreResponse response = receivedResponse; // Wait required to ensure audit logs are flushed into test audit log appender - Thread.sleep(500); + Thread.sleep(1000); // Check message audit logs are produced for each audit log. JsonFormat.Parser protoJSONParser = JsonFormat.parser(); List logJsonObjects = parseMessageJsonLogObjects(testAuditLogAppender.getLogs()); @@ -143,7 +143,7 @@ public void shouldProduceMessageAuditLogsOnError() throws InterruptedException { assertTrue(hasExpectedException); // Wait required to ensure audit logs are flushed into test audit log appender - Thread.sleep(500); + Thread.sleep(1000); List logJsonObjects = parseMessageJsonLogObjects(testAuditLogAppender.getLogs()); assertEquals(logJsonObjects.size(), 1); JsonObject logJsonObject = logJsonObjects.get(0); From 8de5e5675147cd4dfb9121cdf0ca558bd25aaa7e Mon Sep 17 00:00:00 2001 From: Zhu Zhanyan Date: Wed, 19 Aug 2020 20:10:01 +0800 Subject: [PATCH 08/18] Fix to compare the correct lob JsonObject with the right response. --- .../feast/core/logging/CoreLoggingIT.java | 48 +++++++++++-------- 1 file changed, 27 insertions(+), 21 deletions(-) diff --git a/core/src/test/java/feast/core/logging/CoreLoggingIT.java b/core/src/test/java/feast/core/logging/CoreLoggingIT.java index 161c0e257f9..c974ff12f28 100644 --- a/core/src/test/java/feast/core/logging/CoreLoggingIT.java +++ b/core/src/test/java/feast/core/logging/CoreLoggingIT.java @@ -42,6 +42,8 @@ import java.util.List; import java.util.concurrent.ExecutionException; import java.util.stream.Collectors; +import org.apache.beam.vendor.guava.v26_0_jre.com.google.common.collect.Streams; +import org.apache.commons.lang3.tuple.Pair; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.core.LoggerContext; import org.junit.jupiter.api.BeforeAll; @@ -83,13 +85,11 @@ public void shouldProduceMessageAuditLogs() throws InterruptedException { // Generate artifical load on feast core. UpdateStoreRequest request = UpdateStoreRequest.newBuilder().setStore(DataGenerator.getDefaultStore()).build(); - UpdateStoreResponse receivedResponse = null; + List responses = new ArrayList<>(); final int LOAD_SIZE = 200; - final int ASYNC_BURST_SIZE = 5; for (int i = 0; i < LOAD_SIZE; i++) { - receivedResponse = coreService.updateStore(request); + responses.add(coreService.updateStore(request)); } - final UpdateStoreResponse response = receivedResponse; // Wait required to ensure audit logs are flushed into test audit log appender Thread.sleep(1000); @@ -97,25 +97,31 @@ public void shouldProduceMessageAuditLogs() throws InterruptedException { JsonFormat.Parser protoJSONParser = JsonFormat.parser(); List logJsonObjects = parseMessageJsonLogObjects(testAuditLogAppender.getLogs()); - logJsonObjects.forEach( - logObj -> { - // Extract recorded request/response from message logs - try { - String requestJson = logObj.getAsJsonObject("request").toString(); - UpdateStoreRequest.Builder gotRequest = UpdateStoreRequest.newBuilder(); - protoJSONParser.merge(requestJson, gotRequest); + Streams.zip( + logJsonObjects.stream(), + responses.stream(), + (logObj, response) -> Pair.of(logObj, response)) + .forEach( + logObjResponsePair -> { + // Extract recorded request/response from message logs + try { + JsonObject logObj = logObjResponsePair.getLeft(); + UpdateStoreResponse response = logObjResponsePair.getRight(); + String requestJson = logObj.getAsJsonObject("request").toString(); + UpdateStoreRequest.Builder gotRequest = UpdateStoreRequest.newBuilder(); + protoJSONParser.merge(requestJson, gotRequest); - String responseJson = logObj.getAsJsonObject("response").toString(); - UpdateStoreResponse.Builder gotResponse = UpdateStoreResponse.newBuilder(); - protoJSONParser.merge(responseJson, gotResponse); + String responseJson = logObj.getAsJsonObject("response").toString(); + UpdateStoreResponse.Builder gotResponse = UpdateStoreResponse.newBuilder(); + protoJSONParser.merge(responseJson, gotResponse); - // Check that request/response are returned correctlyf - assertThat(gotRequest.build(), equalTo(request)); - assertThat(gotResponse.build(), equalTo(response)); - } catch (InvalidProtocolBufferException e) { - throw new RuntimeException(e); - } - }); + // Check that request/response are returned correctlyf + assertThat(gotRequest.build(), equalTo(request)); + assertThat(gotResponse.build(), equalTo(response)); + } catch (InvalidProtocolBufferException e) { + throw new RuntimeException(e); + } + }); } /** Check that message audit logs are produced when server encounters an error */ From b3a1acd0b63a87672c8c893b63bc1bc1d5e2656f Mon Sep 17 00:00:00 2001 From: Zhu Zhanyan Date: Wed, 19 Aug 2020 20:30:08 +0800 Subject: [PATCH 09/18] Debug response --- core/src/test/java/feast/core/logging/CoreLoggingIT.java | 2 ++ 1 file changed, 2 insertions(+) diff --git a/core/src/test/java/feast/core/logging/CoreLoggingIT.java b/core/src/test/java/feast/core/logging/CoreLoggingIT.java index c974ff12f28..0ce3598d5d8 100644 --- a/core/src/test/java/feast/core/logging/CoreLoggingIT.java +++ b/core/src/test/java/feast/core/logging/CoreLoggingIT.java @@ -109,9 +109,11 @@ public void shouldProduceMessageAuditLogs() throws InterruptedException { UpdateStoreResponse response = logObjResponsePair.getRight(); String requestJson = logObj.getAsJsonObject("request").toString(); UpdateStoreRequest.Builder gotRequest = UpdateStoreRequest.newBuilder(); + System.out.println("got request: " + requestJson); protoJSONParser.merge(requestJson, gotRequest); String responseJson = logObj.getAsJsonObject("response").toString(); + System.out.println("got response: " + responseJson); UpdateStoreResponse.Builder gotResponse = UpdateStoreResponse.newBuilder(); protoJSONParser.merge(responseJson, gotResponse); From e6bb39ba54868254502dd60c3744a9fc209d6d7a Mon Sep 17 00:00:00 2001 From: Zhu Zhanyan Date: Wed, 19 Aug 2020 21:00:27 +0800 Subject: [PATCH 10/18] Reduce load size to make test less flaky --- core/src/test/java/feast/core/logging/CoreLoggingIT.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/core/src/test/java/feast/core/logging/CoreLoggingIT.java b/core/src/test/java/feast/core/logging/CoreLoggingIT.java index 0ce3598d5d8..5529fe0b468 100644 --- a/core/src/test/java/feast/core/logging/CoreLoggingIT.java +++ b/core/src/test/java/feast/core/logging/CoreLoggingIT.java @@ -86,7 +86,7 @@ public void shouldProduceMessageAuditLogs() throws InterruptedException { UpdateStoreRequest request = UpdateStoreRequest.newBuilder().setStore(DataGenerator.getDefaultStore()).build(); List responses = new ArrayList<>(); - final int LOAD_SIZE = 200; + final int LOAD_SIZE = 40; for (int i = 0; i < LOAD_SIZE; i++) { responses.add(coreService.updateStore(request)); } From d6ded00fdab9dfbdf472ca60c2c361b70b813ad4 Mon Sep 17 00:00:00 2001 From: Zhu Zhanyan Date: Wed, 19 Aug 2020 21:29:38 +0800 Subject: [PATCH 11/18] Update test to only check request and response for one call. --- .../feast/core/logging/CoreLoggingIT.java | 51 ++++++------------- 1 file changed, 15 insertions(+), 36 deletions(-) diff --git a/core/src/test/java/feast/core/logging/CoreLoggingIT.java b/core/src/test/java/feast/core/logging/CoreLoggingIT.java index 5529fe0b468..7cfa1388cab 100644 --- a/core/src/test/java/feast/core/logging/CoreLoggingIT.java +++ b/core/src/test/java/feast/core/logging/CoreLoggingIT.java @@ -75,55 +75,34 @@ public static void globalSetUp(@Value("${grpc.server.port}") int coreGrpcPort) coreService.getFeastCoreVersion(GetFeastCoreVersionRequest.getDefaultInstance()); } - /** - * Test Message Audit Logging. Produce artifical load on Feast core and checking the produced - * message logs. - */ + + /** Check that messsage audit log are produced on service call */ @Test - public void shouldProduceMessageAuditLogs() throws InterruptedException { + public void shouldProduceMessageAuditLogsOnCall() throws InterruptedException { testAuditLogAppender.getLogs().clear(); // Generate artifical load on feast core. UpdateStoreRequest request = UpdateStoreRequest.newBuilder().setStore(DataGenerator.getDefaultStore()).build(); - List responses = new ArrayList<>(); - final int LOAD_SIZE = 40; - for (int i = 0; i < LOAD_SIZE; i++) { - responses.add(coreService.updateStore(request)); - } - + UpdateStoreResponse response = coreService.updateStore(request); + // Wait required to ensure audit logs are flushed into test audit log appender Thread.sleep(1000); // Check message audit logs are produced for each audit log. JsonFormat.Parser protoJSONParser = JsonFormat.parser(); List logJsonObjects = parseMessageJsonLogObjects(testAuditLogAppender.getLogs()); + JsonObject logObj = logJsonObjects.get(0); - Streams.zip( - logJsonObjects.stream(), - responses.stream(), - (logObj, response) -> Pair.of(logObj, response)) - .forEach( - logObjResponsePair -> { - // Extract recorded request/response from message logs - try { - JsonObject logObj = logObjResponsePair.getLeft(); - UpdateStoreResponse response = logObjResponsePair.getRight(); - String requestJson = logObj.getAsJsonObject("request").toString(); - UpdateStoreRequest.Builder gotRequest = UpdateStoreRequest.newBuilder(); - System.out.println("got request: " + requestJson); - protoJSONParser.merge(requestJson, gotRequest); + String requestJson = logObj.getAsJsonObject("request").toString(); + UpdateStoreRequest.Builder gotRequest = UpdateStoreRequest.newBuilder(); + protoJSONParser.merge(requestJson, gotRequest); - String responseJson = logObj.getAsJsonObject("response").toString(); - System.out.println("got response: " + responseJson); - UpdateStoreResponse.Builder gotResponse = UpdateStoreResponse.newBuilder(); - protoJSONParser.merge(responseJson, gotResponse); + String responseJson = logObj.getAsJsonObject("response").toString(); + UpdateStoreResponse.Builder gotResponse = UpdateStoreResponse.newBuilder(); + protoJSONParser.merge(responseJson, gotResponse); - // Check that request/response are returned correctlyf - assertThat(gotRequest.build(), equalTo(request)); - assertThat(gotResponse.build(), equalTo(response)); - } catch (InvalidProtocolBufferException e) { - throw new RuntimeException(e); - } - }); + // Check that request/response are returned correctlyf + assertThat(gotRequest.build(), equalTo(request)); + assertThat(gotResponse.build(), equalTo(response)); } /** Check that message audit logs are produced when server encounters an error */ From 74b07bceeafde1bdb62c1d137d5337c6f9bea752 Mon Sep 17 00:00:00 2001 From: Zhu Zhanyan Date: Thu, 20 Aug 2020 09:13:30 +0800 Subject: [PATCH 12/18] Fix compile failure due to uncaught exception. --- core/src/test/java/feast/core/logging/CoreLoggingIT.java | 8 +++----- 1 file changed, 3 insertions(+), 5 deletions(-) diff --git a/core/src/test/java/feast/core/logging/CoreLoggingIT.java b/core/src/test/java/feast/core/logging/CoreLoggingIT.java index 7cfa1388cab..d6c0ed4f5da 100644 --- a/core/src/test/java/feast/core/logging/CoreLoggingIT.java +++ b/core/src/test/java/feast/core/logging/CoreLoggingIT.java @@ -42,8 +42,6 @@ import java.util.List; import java.util.concurrent.ExecutionException; import java.util.stream.Collectors; -import org.apache.beam.vendor.guava.v26_0_jre.com.google.common.collect.Streams; -import org.apache.commons.lang3.tuple.Pair; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.core.LoggerContext; import org.junit.jupiter.api.BeforeAll; @@ -75,16 +73,16 @@ public static void globalSetUp(@Value("${grpc.server.port}") int coreGrpcPort) coreService.getFeastCoreVersion(GetFeastCoreVersionRequest.getDefaultInstance()); } - /** Check that messsage audit log are produced on service call */ @Test - public void shouldProduceMessageAuditLogsOnCall() throws InterruptedException { + public void shouldProduceMessageAuditLogsOnCall() + throws InterruptedException, InvalidProtocolBufferException { testAuditLogAppender.getLogs().clear(); // Generate artifical load on feast core. UpdateStoreRequest request = UpdateStoreRequest.newBuilder().setStore(DataGenerator.getDefaultStore()).build(); UpdateStoreResponse response = coreService.updateStore(request); - + // Wait required to ensure audit logs are flushed into test audit log appender Thread.sleep(1000); // Check message audit logs are produced for each audit log. From 297da1f6842bccba270242a4116478376ac36c03 Mon Sep 17 00:00:00 2001 From: Zhu Zhanyan Date: Thu, 20 Aug 2020 10:14:37 +0800 Subject: [PATCH 13/18] Add method name filter to prevent logs from tests from interfering with each other. --- .../feast/core/logging/CoreLoggingIT.java | 27 ++++++++++++------- 1 file changed, 17 insertions(+), 10 deletions(-) diff --git a/core/src/test/java/feast/core/logging/CoreLoggingIT.java b/core/src/test/java/feast/core/logging/CoreLoggingIT.java index d6c0ed4f5da..68b297b5c9f 100644 --- a/core/src/test/java/feast/core/logging/CoreLoggingIT.java +++ b/core/src/test/java/feast/core/logging/CoreLoggingIT.java @@ -38,7 +38,6 @@ import io.grpc.ManagedChannelBuilder; import io.grpc.Status.Code; import io.grpc.StatusRuntimeException; -import java.util.ArrayList; import java.util.List; import java.util.concurrent.ExecutionException; import java.util.stream.Collectors; @@ -77,7 +76,6 @@ public static void globalSetUp(@Value("${grpc.server.port}") int coreGrpcPort) @Test public void shouldProduceMessageAuditLogsOnCall() throws InterruptedException, InvalidProtocolBufferException { - testAuditLogAppender.getLogs().clear(); // Generate artifical load on feast core. UpdateStoreRequest request = UpdateStoreRequest.newBuilder().setStore(DataGenerator.getDefaultStore()).build(); @@ -87,9 +85,15 @@ public void shouldProduceMessageAuditLogsOnCall() Thread.sleep(1000); // Check message audit logs are produced for each audit log. JsonFormat.Parser protoJSONParser = JsonFormat.parser(); - List logJsonObjects = parseMessageJsonLogObjects(testAuditLogAppender.getLogs()); + // filter by method name to ensure logs from other tests do not interfere with test + List logJsonObjects = + parseMessageJsonLogObjects(testAuditLogAppender.getLogs()).stream() + .filter(logObj -> logObj.get("method").getAsString().equals("UpdateStore")) + .collect(Collectors.toList()); + assertEquals(1, logJsonObjects.size()); JsonObject logObj = logJsonObjects.get(0); + // Extract & Check that request/response are returned correctly String requestJson = logObj.getAsJsonObject("request").toString(); UpdateStoreRequest.Builder gotRequest = UpdateStoreRequest.newBuilder(); protoJSONParser.merge(requestJson, gotRequest); @@ -98,7 +102,6 @@ public void shouldProduceMessageAuditLogsOnCall() UpdateStoreResponse.Builder gotResponse = UpdateStoreResponse.newBuilder(); protoJSONParser.merge(responseJson, gotResponse); - // Check that request/response are returned correctlyf assertThat(gotRequest.build(), equalTo(request)); assertThat(gotResponse.build(), equalTo(response)); } @@ -106,8 +109,7 @@ public void shouldProduceMessageAuditLogsOnCall() /** Check that message audit logs are produced when server encounters an error */ @Test public void shouldProduceMessageAuditLogsOnError() throws InterruptedException { - testAuditLogAppender.getLogs().clear(); - // send a bad request which should cause Core to error + // Send a bad request which should cause Core to error ListFeatureSetsRequest request = ListFeatureSetsRequest.newBuilder() .setFilter( @@ -129,9 +131,15 @@ public void shouldProduceMessageAuditLogsOnError() throws InterruptedException { // Wait required to ensure audit logs are flushed into test audit log appender Thread.sleep(1000); - List logJsonObjects = parseMessageJsonLogObjects(testAuditLogAppender.getLogs()); - assertEquals(logJsonObjects.size(), 1); + // filter by method name to ensure logs from other tests do not interfere with test + List logJsonObjects = + parseMessageJsonLogObjects(testAuditLogAppender.getLogs()).stream() + .filter(logObj -> logObj.get("method").getAsString().equals("ListFeatureSets")) + .collect(Collectors.toList()); + + assertEquals(1, logJsonObjects.size()); JsonObject logJsonObject = logJsonObjects.get(0); + // Check correct status code is tracked on error. assertEquals(logJsonObject.get("statusCode").getAsString(), statusCode.toString()); } @@ -139,8 +147,7 @@ public void shouldProduceMessageAuditLogsOnError() throws InterruptedException { private List parseMessageJsonLogObjects(List logsStrings) { JsonParser jsonParser = new JsonParser(); // copy to prevent concurrent modification. - List logsStringsCopy = new ArrayList(logsStrings); - return logsStringsCopy.stream() + return logsStrings.stream() .map(logJSON -> jsonParser.parse(logJSON).getAsJsonObject()) // Filter to only include message audit logs .filter( From 03b7eb18a2de6909f80a9fd4989bddf77d1e011e Mon Sep 17 00:00:00 2001 From: Zhu Zhanyan Date: Thu, 20 Aug 2020 14:28:28 +0800 Subject: [PATCH 14/18] Add intergration test to check that message logs are produced correctly under load. --- .../feast/core/logging/CoreLoggingIT.java | 70 +++++++++++++++---- 1 file changed, 55 insertions(+), 15 deletions(-) diff --git a/core/src/test/java/feast/core/logging/CoreLoggingIT.java b/core/src/test/java/feast/core/logging/CoreLoggingIT.java index 68b297b5c9f..db78ff51c0f 100644 --- a/core/src/test/java/feast/core/logging/CoreLoggingIT.java +++ b/core/src/test/java/feast/core/logging/CoreLoggingIT.java @@ -21,6 +21,8 @@ import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertTrue; +import com.google.common.util.concurrent.Futures; +import com.google.common.util.concurrent.ListenableFuture; import com.google.gson.JsonObject; import com.google.gson.JsonParser; import com.google.protobuf.InvalidProtocolBufferException; @@ -30,14 +32,18 @@ import feast.core.it.DataGenerator; import feast.proto.core.CoreServiceGrpc; import feast.proto.core.CoreServiceGrpc.CoreServiceBlockingStub; +import feast.proto.core.CoreServiceGrpc.CoreServiceFutureStub; import feast.proto.core.CoreServiceProto.GetFeastCoreVersionRequest; import feast.proto.core.CoreServiceProto.ListFeatureSetsRequest; +import feast.proto.core.CoreServiceProto.ListStoresRequest; +import feast.proto.core.CoreServiceProto.ListStoresResponse; import feast.proto.core.CoreServiceProto.UpdateStoreRequest; import feast.proto.core.CoreServiceProto.UpdateStoreResponse; import io.grpc.Channel; import io.grpc.ManagedChannelBuilder; import io.grpc.Status.Code; import io.grpc.StatusRuntimeException; +import java.util.LinkedList; import java.util.List; import java.util.concurrent.ExecutionException; import java.util.stream.Collectors; @@ -56,20 +62,26 @@ public class CoreLoggingIT extends BaseIT { private static TestLogAppender testAuditLogAppender; private static CoreServiceBlockingStub coreService; + private static CoreServiceFutureStub asyncCoreService; @BeforeAll public static void globalSetUp(@Value("${grpc.server.port}") int coreGrpcPort) throws InterruptedException, ExecutionException { LoggerContext logContext = (LoggerContext) LogManager.getContext(false); + // NOTE: As log appender state is shared across tests use a different method + // for each test and filter by method name to ensure that you only get logs + // for a specific test. testAuditLogAppender = logContext.getConfiguration().getAppender("TestAuditLogAppender"); - // Connect to core. + // Connect to core service. Channel channel = ManagedChannelBuilder.forAddress("localhost", coreGrpcPort).usePlaintext().build(); coreService = CoreServiceGrpc.newBlockingStub(channel); + asyncCoreService = CoreServiceGrpc.newFutureStub(channel); - // Preflight a request to core service stub to verify connection + // Preflight a request to core service stubs to verify connection coreService.getFeastCoreVersion(GetFeastCoreVersionRequest.getDefaultInstance()); + asyncCoreService.getFeastCoreVersion(GetFeastCoreVersionRequest.getDefaultInstance()).get(); } /** Check that messsage audit log are produced on service call */ @@ -85,11 +97,9 @@ public void shouldProduceMessageAuditLogsOnCall() Thread.sleep(1000); // Check message audit logs are produced for each audit log. JsonFormat.Parser protoJSONParser = JsonFormat.parser(); - // filter by method name to ensure logs from other tests do not interfere with test + // Pull message audit logs logs from test log appender List logJsonObjects = - parseMessageJsonLogObjects(testAuditLogAppender.getLogs()).stream() - .filter(logObj -> logObj.get("method").getAsString().equals("UpdateStore")) - .collect(Collectors.toList()); + parseMessageJsonLogObjects(testAuditLogAppender.getLogs(), "UpdateStore"); assertEquals(1, logJsonObjects.size()); JsonObject logObj = logJsonObjects.get(0); @@ -131,11 +141,9 @@ public void shouldProduceMessageAuditLogsOnError() throws InterruptedException { // Wait required to ensure audit logs are flushed into test audit log appender Thread.sleep(1000); - // filter by method name to ensure logs from other tests do not interfere with test + // Pull message audit logs logs from test log appender List logJsonObjects = - parseMessageJsonLogObjects(testAuditLogAppender.getLogs()).stream() - .filter(logObj -> logObj.get("method").getAsString().equals("ListFeatureSets")) - .collect(Collectors.toList()); + parseMessageJsonLogObjects(testAuditLogAppender.getLogs(), "ListFeatureSets"); assertEquals(1, logJsonObjects.size()); JsonObject logJsonObject = logJsonObjects.get(0); @@ -143,8 +151,37 @@ public void shouldProduceMessageAuditLogsOnError() throws InterruptedException { assertEquals(logJsonObject.get("statusCode").getAsString(), statusCode.toString()); } - /** Filter and Parse out Message Audit Logs from the given logsStrings */ - private List parseMessageJsonLogObjects(List logsStrings) { + /** Check that expected message audit logs are produced when under load. */ + @Test + public void shouldProduceExpectedNumberOfAuditLogs() + throws InterruptedException, ExecutionException { + // Generate artifical requests on core to simulate load. + int LOAD_SIZE = 40; // Total number of requests to send. + int BURST_SIZE = 5; // Number of requests to send at once. + + List responses = new LinkedList<>(); + for (int i = 0; i < LOAD_SIZE; i += 5) { + List> futures = new LinkedList<>(); + for (int j = 0; j < BURST_SIZE; j++) { + futures.add(asyncCoreService.listStores(ListStoresRequest.getDefaultInstance())); + } + + responses.addAll(Futures.allAsList(futures).get()); + } + // Wait required to ensure audit logs are flushed into test audit log appender + Thread.sleep(1000); + + // Pull message audit logs logs from test log appender + List logJsonObjects = + parseMessageJsonLogObjects(testAuditLogAppender.getLogs(), "ListStores"); + + assertEquals(responses.size(), logJsonObjects.size()); + } + + /** + * Filter and Parse out Message Audit Logs from the given logsStrings for the given method name + */ + private List parseMessageJsonLogObjects(List logsStrings, String methodName) { JsonParser jsonParser = new JsonParser(); // copy to prevent concurrent modification. return logsStrings.stream() @@ -153,9 +190,12 @@ private List parseMessageJsonLogObjects(List logsStrings) { .filter( logObj -> logObj - .getAsJsonPrimitive("kind") - .getAsString() - .equals(AuditLogEntryKind.MESSAGE.toString())) + .getAsJsonPrimitive("kind") + .getAsString() + .equals(AuditLogEntryKind.MESSAGE.toString()) + // filter by method name to ensure logs from other tests do not interfere with + // test + && logObj.get("method").getAsString().equals(methodName)) .collect(Collectors.toList()); } } From e1b173f89d41b407cd1d84f97dc05d9d40316f72 Mon Sep 17 00:00:00 2001 From: Zhu Zhanyan Date: Sat, 22 Aug 2020 12:23:09 +0800 Subject: [PATCH 15/18] Fix imports and log4j2 not able to find config file --- core/src/test/java/feast/core/logging/CoreLoggingIT.java | 4 ++-- core/src/test/resources/{log4j2-test.xml => log4j2.xml} | 0 2 files changed, 2 insertions(+), 2 deletions(-) rename core/src/test/resources/{log4j2-test.xml => log4j2.xml} (100%) diff --git a/core/src/test/java/feast/core/logging/CoreLoggingIT.java b/core/src/test/java/feast/core/logging/CoreLoggingIT.java index db78ff51c0f..80e687f568f 100644 --- a/core/src/test/java/feast/core/logging/CoreLoggingIT.java +++ b/core/src/test/java/feast/core/logging/CoreLoggingIT.java @@ -27,9 +27,9 @@ import com.google.gson.JsonParser; import com.google.protobuf.InvalidProtocolBufferException; import com.google.protobuf.util.JsonFormat; +import feast.common.it.BaseIT; +import feast.common.it.DataGenerator; import feast.common.logging.entry.AuditLogEntryKind; -import feast.core.it.BaseIT; -import feast.core.it.DataGenerator; import feast.proto.core.CoreServiceGrpc; import feast.proto.core.CoreServiceGrpc.CoreServiceBlockingStub; import feast.proto.core.CoreServiceGrpc.CoreServiceFutureStub; diff --git a/core/src/test/resources/log4j2-test.xml b/core/src/test/resources/log4j2.xml similarity index 100% rename from core/src/test/resources/log4j2-test.xml rename to core/src/test/resources/log4j2.xml From 81135c2ee0c40e64dcd8b9e0a782de587ab3a21c Mon Sep 17 00:00:00 2001 From: Zhu Zhanyan Date: Wed, 26 Aug 2020 14:22:00 +0800 Subject: [PATCH 16/18] Fix issue with CoreLoggingIT TestLogAppender being null due to class not found by log4j2 --- core/src/test/java/feast/core/logging/TestLogAppender.java | 7 ++++++- core/src/test/resources/log4j2.xml | 2 +- 2 files changed, 7 insertions(+), 2 deletions(-) diff --git a/core/src/test/java/feast/core/logging/TestLogAppender.java b/core/src/test/java/feast/core/logging/TestLogAppender.java index 21c5bbd69d0..b1908c0db2c 100644 --- a/core/src/test/java/feast/core/logging/TestLogAppender.java +++ b/core/src/test/java/feast/core/logging/TestLogAppender.java @@ -20,6 +20,8 @@ import java.util.ArrayList; import java.util.List; import lombok.Getter; +import org.apache.logging.log4j.core.Appender; +import org.apache.logging.log4j.core.Core; import org.apache.logging.log4j.core.Filter; import org.apache.logging.log4j.core.Layout; import org.apache.logging.log4j.core.LogEvent; @@ -32,7 +34,10 @@ import org.apache.logging.log4j.core.layout.PatternLayout; /** Test Log Appender used for collecting logs for testing logging. */ -@Plugin(name = "TestLogAppender", category = "Core", elementType = "appender", printObject = true) +@Plugin( + name = "TestLogAppender", + category = Core.CATEGORY_NAME, + elementType = Appender.ELEMENT_TYPE) @Getter public class TestLogAppender extends AbstractAppender { private List logs; diff --git a/core/src/test/resources/log4j2.xml b/core/src/test/resources/log4j2.xml index 90945373bfe..9f785952e0c 100644 --- a/core/src/test/resources/log4j2.xml +++ b/core/src/test/resources/log4j2.xml @@ -16,7 +16,7 @@ ~ --> - + %d{yyyy-MM-dd HH:mm:ss.SSS} %5p ${hostName} --- [%15.15t] %-40.40c{1.} : %m%n%ex From 794e8c4b0855a162b0268413ef64403d8404ab7e Mon Sep 17 00:00:00 2001 From: Zhu Zhanyan Date: Tue, 1 Sep 2020 12:35:23 +0800 Subject: [PATCH 17/18] Remove unused getters in MessageAuditLogEntry. --- .../logging/entry/MessageAuditLogEntry.java | 22 +------------------ 1 file changed, 1 insertion(+), 21 deletions(-) diff --git a/common/src/main/java/feast/common/logging/entry/MessageAuditLogEntry.java b/common/src/main/java/feast/common/logging/entry/MessageAuditLogEntry.java index f158dae4221..745cc1283ae 100644 --- a/common/src/main/java/feast/common/logging/entry/MessageAuditLogEntry.java +++ b/common/src/main/java/feast/common/logging/entry/MessageAuditLogEntry.java @@ -66,44 +66,24 @@ public abstract class MessageAuditLogEntry extends AuditLogEntry { public abstract static class Builder { public abstract Builder setId(UUID id); - public abstract UUID getId(); - public abstract Builder setComponent(String component); - public abstract String getComponent(); - - public abstract Builder setVersion(String version); - - public abstract String getVersion(); + public abstract Builder setVersion(String component); public abstract Builder setKind(AuditLogEntryKind kind); - public abstract AuditLogEntryKind getKind(); - public abstract Builder setService(String name); - public abstract String getService(); - public abstract Builder setMethod(String name); - public abstract String getMethod(); - public abstract Builder setRequest(Message request); - public abstract Message getRequest(); - public abstract Builder setResponse(Message response); - public abstract Message getResponse(); - public abstract Builder setIdentity(String identity); - public abstract String getIdentity(); - public abstract Builder setStatusCode(Code statusCode); - public abstract Code getStatusCode(); - public abstract MessageAuditLogEntry build(); } From fc75f83ef8f149651256181403f95bbfa6120783 Mon Sep 17 00:00:00 2001 From: Zhu Zhanyan Date: Thu, 3 Sep 2020 10:29:47 +0800 Subject: [PATCH 18/18] Update CoreLoggingIT test to check that expected contents of logs produced under load. --- .../feast/core/logging/CoreLoggingIT.java | 39 +++++++++++++++++-- 1 file changed, 35 insertions(+), 4 deletions(-) diff --git a/core/src/test/java/feast/core/logging/CoreLoggingIT.java b/core/src/test/java/feast/core/logging/CoreLoggingIT.java index 80e687f568f..53e55ebe695 100644 --- a/core/src/test/java/feast/core/logging/CoreLoggingIT.java +++ b/core/src/test/java/feast/core/logging/CoreLoggingIT.java @@ -21,6 +21,7 @@ import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertTrue; +import com.google.common.collect.Streams; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.ListenableFuture; import com.google.gson.JsonObject; @@ -47,6 +48,7 @@ import java.util.List; import java.util.concurrent.ExecutionException; import java.util.stream.Collectors; +import org.apache.commons.lang3.tuple.Pair; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.core.LoggerContext; import org.junit.jupiter.api.BeforeAll; @@ -153,17 +155,18 @@ public void shouldProduceMessageAuditLogsOnError() throws InterruptedException { /** Check that expected message audit logs are produced when under load. */ @Test - public void shouldProduceExpectedNumberOfAuditLogs() + public void shouldProduceExpectedAuditLogsUnderLoad() throws InterruptedException, ExecutionException { // Generate artifical requests on core to simulate load. int LOAD_SIZE = 40; // Total number of requests to send. int BURST_SIZE = 5; // Number of requests to send at once. + ListStoresRequest request = ListStoresRequest.getDefaultInstance(); List responses = new LinkedList<>(); for (int i = 0; i < LOAD_SIZE; i += 5) { List> futures = new LinkedList<>(); for (int j = 0; j < BURST_SIZE; j++) { - futures.add(asyncCoreService.listStores(ListStoresRequest.getDefaultInstance())); + futures.add(asyncCoreService.listStores(request)); } responses.addAll(Futures.allAsList(futures).get()); @@ -171,11 +174,39 @@ public void shouldProduceExpectedNumberOfAuditLogs() // Wait required to ensure audit logs are flushed into test audit log appender Thread.sleep(1000); - // Pull message audit logs logs from test log appender + // Pull message audit logs from test log appender List logJsonObjects = parseMessageJsonLogObjects(testAuditLogAppender.getLogs(), "ListStores"); - assertEquals(responses.size(), logJsonObjects.size()); + + // Extract & Check that request/response are returned correctly + JsonFormat.Parser protoJSONParser = JsonFormat.parser(); + Streams.zip( + responses.stream(), + logJsonObjects.stream(), + (response, logObj) -> Pair.of(response, logObj)) + .forEach( + responseLogJsonPair -> { + ListStoresResponse response = responseLogJsonPair.getLeft(); + JsonObject logObj = responseLogJsonPair.getRight(); + + ListStoresRequest.Builder gotRequest = null; + ListStoresResponse.Builder gotResponse = null; + try { + String requestJson = logObj.getAsJsonObject("request").toString(); + gotRequest = ListStoresRequest.newBuilder(); + protoJSONParser.merge(requestJson, gotRequest); + + String responseJson = logObj.getAsJsonObject("response").toString(); + gotResponse = ListStoresResponse.newBuilder(); + protoJSONParser.merge(responseJson, gotResponse); + } catch (InvalidProtocolBufferException e) { + throw new RuntimeException(e); + } + + assertThat(gotRequest.build(), equalTo(request)); + assertThat(gotResponse.build(), equalTo(response)); + }); } /**