diff --git a/zookeeper-docs/src/main/resources/markdown/zookeeperAuditLogs.md b/zookeeper-docs/src/main/resources/markdown/zookeeperAuditLogs.md
index b3954ad676d..4e03b13466d 100644
--- a/zookeeper-docs/src/main/resources/markdown/zookeeperAuditLogs.md
+++ b/zookeeper-docs/src/main/resources/markdown/zookeeperAuditLogs.md
@@ -17,7 +17,8 @@ limitations under the License.
# ZooKeeper Audit Logging
* [ZooKeeper Audit Logs](#ch_auditLogs)
-* [ZooKeeper Audit Log Configuration](#ch_reconfig_format)
+* [ZooKeeper Audit Log Configuration](#ch_auditConfig)
+* [Enhanced audit metadata (schema v2)](#ch_auditV2)
* [Who is taken as user in audit logs?](#ch_zkAuditUser)
@@ -36,9 +37,9 @@ The audit log captures detailed information for the operations that are selected
|session | client session id |
|user | comma separated list of users who are associate with a client session. For more on this, see [Who is taken as user in audit logs](#ch_zkAuditUser).
|ip | client IP address
-|operation | any one of the selected operations for audit. Possible values are(serverStart, serverStop, create, delete, setData, setAcl, multiOperation, reconfig, ephemeralZNodeDeleteOnSessionClose)
+|operation | any one of the selected operations for audit. Possible values are(serverStart, serverStop, create, delete, setData, setAcl, multiOperation, reconfig, ephemeralZNodeDeletionOnSessionCloseOrExpire)
|znode | path of the znode
-|znode type | type of znode in case of creation operation
+|znode_type | type of znode in case of creation operation
|acl | String representation of znode ACL like cdrwa(create, delete,read, write, admin). This is logged only for setAcl operation
|result | result of the operation. Possible values are (success/failure/invoked). Result "invoked" is used for serverStop operation because stop is logged before ensuring that server actually stopped.
@@ -99,6 +100,145 @@ Audit logging is done using logback. Following is the default logback configurat
Change above configuration to customize the auditlog file, number of backups, max file size, custom audit logger etc.
+
+
+## Enhanced audit metadata (schema v2)
+
+The optional Java system property `zookeeper.audit.enhanced.enable` defaults to
+`false`. Set `-Dzookeeper.audit.enhanced.enable=true` in addition to enabling
+`zookeeper.audit.enable` to emit v2 records. Enhanced logging does not enable
+audit logging by itself. With enhanced logging off, ordinary legacy formatting
+and the existing logging APIs are preserved. TTL creates are audited in both
+modes, including sequential TTL creates.
+
+V2 keeps the existing field names and `result=success/failure/invoked`, and adds:
+
+| Key | Meaning |
+| --- | --- |
+| `schema_version` | `2`. Its absence identifies legacy output. |
+| `data_length` | Attempted payload size in bytes for create variants and `setData`, not request size, character count, or resulting znode size. Null and empty payloads are `0`; unavailable or undecodable payloads omit this field. A failed attempt can have a nonzero length. |
+| `error_code` | Known numeric ZooKeeper code, such as `0` (OK), `-101` (NONODE), `-102` (NOAUTH), `-103` (BADVERSION), or `-110` (NODEEXISTS). Omitted when unavailable. |
+| `outcome` | `committed`, `failed`, `rolled_back`, or `unknown`, as described below. |
+| `cxid` | Known client request ID in decimal. All members of a multi share this ID. |
+| `zxid` | Known transaction ID in decimal, not the server's latest zxid at reply time. A failed transaction may also have a zxid. |
+| `multi_index` | Zero-based position in the **complete** multi request, including non-mutating checks. Omitted on single operations and multi parent records. |
+
+No znode payloads are included. For enhanced `setAcl` records, `acl` retains
+schemes and permissions, but identities are conservative: `world:anyone` is
+retained, valid built-in digest identities use the provider's username extraction
+(not the digest), and other identities are `[redacted]`. This also applies to
+failed ACL attempts, where identities might be malformed or contain credentials.
+Unknown custom identity representations are not printed as ACL identities.
+
+For server-generated v2 records, the `user` field also uses a conservative
+provider policy rather than trusting a custom provider's default `getUserName`.
+Only the concrete built-in digest, IP, SASL and X509 providers are trusted, and
+their identity syntax must validate before their username extraction is used.
+Digest credentials are reduced to usernames; X509 identities are distinguished
+names, not certificate bodies. Unknown or malformed identities, custom providers
+and provider subclasses are represented as `[redacted]`. Merely overriding
+`getUserName` does not opt a custom provider into this trust policy. Legacy user
+extraction is unchanged when enhanced logging is off.
+
+Custom callers of the logging APIs must supply sanitized user and ACL metadata;
+the string-based APIs cannot infer credentials from arbitrary strings.
+
+### Results and transaction outcomes
+
+* `committed` means the available transaction result reports an applied write.
+ These records have `result=success` and `error_code=0`. It does not guarantee
+ that the reply reached the client or the log reached durable storage.
+* `failed` means the operation or whole multi failed according to its known
+ result. For rejected single transactions, the request exception takes
+ precedence when the reply path uses it instead of the transaction error.
+* `rolled_back` means a mutation was not applied because its atomic multi
+ failed. This includes prepared members discarded by rollback and members
+ skipped after the failure. They have `result=failure`, **even when their
+ individual `error_code` is `0`**. Later skipped members normally have code
+ `-2` (RUNTIMEINCONSISTENCY); a first failure with that code is still `failed`.
+* `unknown` means the available metadata cannot establish an operation's
+ outcome. Missing or mismatched multi results are not treated as commits,
+ and untrustworthy per-member error codes are omitted. If the whole multi
+ is known to have failed, these members still have `result=failure`;
+ otherwise unknown operations use `result=invoked`.
+
+Failed multis keep their `multiOperation` parent failure record, followed in v2
+by records for the attempted mutations when the request is decodable. The parent
+is retained even if request decoding fails. Successful multis emit individual
+mutation records, not a parent success record. Request and result members are
+paired by position, so repeated paths cannot overwrite each other's metadata.
+Successful sequential creates use the final path returned by the transaction;
+failed or rolled-back creates use the attempted path.
+
+### Supported operations and identity
+
+The audited operations remain creates (`create`, `create2`, container and TTL
+variants), deletes (including container deletes), `setData`, `setAcl`, write
+multis, reconfiguration, server start/stop, and existing system ephemeral-node
+deletions. Ordinary reads, read multis and checks do not produce audit records.
+Checks inside a write multi still occupy an index.
+
+Existing lifecycle events and callers of the legacy logging overloads receive
+`outcome=unknown` in v2 when they do not supply transaction metadata; their
+existing `result` is unchanged. No session/authentication binding events are
+added by this schema.
+
+System ephemeral-node deletion records preserve the **server actor**, the
+affected session and path, and omit client IP. V2 adds the deletion transaction's
+zxid, known error code and outcome. A client request ID is not available at this
+deletion hook, so it is omitted. These system records can be emitted on multiple
+replicas or during replay. Downstream consumers can deduplicate using ensemble
+identity plus `zxid`, `operation`, `session` and `znode`; client multi mutations
+also require `multi_index` to distinguish repeated paths.
+
+### Escaping and parsing
+
+Fields are separated by literal tabs. Split each field at its **first** equals
+sign; equals signs inside values are unchanged. In v2 only, value formatting
+reversibly escapes backslash as `\\`, tab as `\t`, carriage return as `\r`, and
+newline as `\n`. A literal backslash followed by `t` is therefore `\\t`, distinct
+from an escaped tab. Decode left to right, consuming one escape pair at a time;
+do not use successive global replacements to unescape. For example, the user
+value containing a tab and newline is rendered as `user=team\tname\n` on a single
+physical log line. Enum-derived keys and values are independent of the server's
+default locale. Custom `AuditLogger` implementations receive raw event values;
+`AuditEvent.toString()` is the production text-formatting path.
+
+### Mixed versions, rollout, rollback and coverage
+
+Consumers should recognize v2 per record using `schema_version=2`, tolerate
+unknown additive fields, and accept omitted optional fields. Do not apply v2
+unescaping to legacy records. Prepare consumers for mixed legacy/v2 output before
+enabling the property on servers, then roll it out gradually. Roll back emission
+by removing the enhanced property or setting it to `false` on restart. Each
+audited request captures the enhanced mode before metadata extraction, and uses
+that same decision for user/ACL sanitization, schema construction and every
+parent/member record of a multi. An in-flight request keeps its captured mode if
+the property changes; a subsequent request captures the new value. Direct
+provider logging captures its mode per event. The base audit enablement retains
+its startup behavior. No wire, persistence, authentication, ACL, quota or payload
+limit changes are required.
+
+The bounded-cardinality server counter `audit_errors` counts detected audit
+metadata extraction/correlation failures and runtime failures reported by the
+audit logger. Such failures are logged without rejecting an otherwise valid
+client operation. A metadata error can still yield an audit record with omitted
+fields, so this counter is **not** a dropped-record count. It has no per-path,
+per-user or per-session labels.
+The counter update and diagnostic error logging are independently best-effort:
+runtime failures in either reporting backend are isolated without retries, so
+they cannot reject an applied write or interrupt system deletion. A failing
+metrics backend can therefore also leave detected audit errors uncounted.
+
+Audit coverage is limited to the existing transaction-audit hooks. Rejections
+before a hook, connection-level failures and later reply/send failures need
+separate instrumentation; these records are not a complete request-delivery
+ledger. Log delivery is also a separate concern: an asynchronous appender using
+`neverBlock=true` can discard records without reporting an exception.
+`audit_errors` cannot measure those silent drops, filtering, retention loss, or
+downstream delivery gaps. Validate source-log coverage and delivery independently
+before relying on the stream for complete accounting.
+
## Who is taken as user in audit logs?
diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/audit/AuditConstants.java b/zookeeper-server/src/main/java/org/apache/zookeeper/audit/AuditConstants.java
index 22fd8567abe..1bfe96b363b 100644
--- a/zookeeper-server/src/main/java/org/apache/zookeeper/audit/AuditConstants.java
+++ b/zookeeper-server/src/main/java/org/apache/zookeeper/audit/AuditConstants.java
@@ -18,6 +18,8 @@
package org.apache.zookeeper.audit;
public final class AuditConstants {
+ public static final String SCHEMA_VERSION = "2";
+
private AuditConstants() {
//Utility classes should not have public constructors
}
diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/audit/AuditEvent.java b/zookeeper-server/src/main/java/org/apache/zookeeper/audit/AuditEvent.java
index e499552a948..bd3a7842f95 100644
--- a/zookeeper-server/src/main/java/org/apache/zookeeper/audit/AuditEvent.java
+++ b/zookeeper-server/src/main/java/org/apache/zookeeper/audit/AuditEvent.java
@@ -18,6 +18,7 @@
package org.apache.zookeeper.audit;
import java.util.LinkedHashMap;
+import java.util.Locale;
import java.util.Map;
import java.util.Set;
@@ -43,12 +44,12 @@ public Set> getLogEntries() {
void addEntry(FieldName fieldName, String value) {
if (value != null) {
- logEntries.put(fieldName.name().toLowerCase(), value);
+ logEntries.put(fieldName.name().toLowerCase(Locale.ROOT), value);
}
}
public String getValue(FieldName fieldName) {
- return logEntries.get(fieldName.name().toLowerCase());
+ return logEntries.get(fieldName.name().toLowerCase(Locale.ROOT));
}
public Result getResult() {
@@ -63,6 +64,7 @@ public Result getResult() {
@Override
public String toString() {
StringBuilder buffer = new StringBuilder();
+ boolean enhanced = AuditConstants.SCHEMA_VERSION.equals(getValue(FieldName.SCHEMA_VERSION));
boolean first = true;
for (Map.Entry entry : logEntries.entrySet()) {
String key = entry.getKey();
@@ -75,7 +77,7 @@ public String toString() {
buffer.append(PAIR_SEPARATOR);
}
buffer.append(key).append(KEY_VAL_SEPARATOR)
- .append(value);
+ .append(enhanced ? escape(value) : value);
}
}
//add result field
@@ -83,16 +85,27 @@ public String toString() {
buffer.append(PAIR_SEPARATOR);
}
buffer.append("result").append(KEY_VAL_SEPARATOR)
- .append(result.name().toLowerCase());
+ .append(result.name().toLowerCase(Locale.ROOT));
return buffer.toString();
}
+ private static String escape(String value) {
+ return value.replace("\\", "\\\\")
+ .replace("\t", "\\t")
+ .replace("\r", "\\r")
+ .replace("\n", "\\n");
+ }
+
public enum FieldName {
- USER, OPERATION, IP, ACL, ZNODE, SESSION, ZNODE_TYPE
+ USER, OPERATION, IP, ACL, ZNODE, SESSION, ZNODE_TYPE,
+ SCHEMA_VERSION, DATA_LENGTH, ERROR_CODE, OUTCOME, CXID, ZXID, MULTI_INDEX
}
public enum Result {
SUCCESS, FAILURE, INVOKED
}
-}
+ public enum Outcome {
+ COMMITTED, FAILED, ROLLED_BACK, UNKNOWN
+ }
+}
diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/audit/AuditHelper.java b/zookeeper-server/src/main/java/org/apache/zookeeper/audit/AuditHelper.java
index ce0d58a87ba..4731cbe58a3 100644
--- a/zookeeper-server/src/main/java/org/apache/zookeeper/audit/AuditHelper.java
+++ b/zookeeper-server/src/main/java/org/apache/zookeeper/audit/AuditHelper.java
@@ -18,23 +18,35 @@
package org.apache.zookeeper.audit;
import java.io.IOException;
-import java.util.HashMap;
-import java.util.Map;
+import java.nio.ByteBuffer;
+import java.util.List;
+import java.util.Locale;
import org.apache.jute.Record;
import org.apache.zookeeper.CreateMode;
import org.apache.zookeeper.KeeperException;
+import org.apache.zookeeper.KeeperException.Code;
import org.apache.zookeeper.MultiOperationRecord;
import org.apache.zookeeper.Op;
import org.apache.zookeeper.ZKUtil;
-import org.apache.zookeeper.ZooDefs;
+import org.apache.zookeeper.ZooDefs.OpCode;
+import org.apache.zookeeper.audit.AuditEvent.Outcome;
import org.apache.zookeeper.audit.AuditEvent.Result;
+import org.apache.zookeeper.data.ACL;
+import org.apache.zookeeper.data.Id;
import org.apache.zookeeper.proto.CreateRequest;
+import org.apache.zookeeper.proto.CreateTTLRequest;
import org.apache.zookeeper.proto.DeleteRequest;
import org.apache.zookeeper.proto.SetACLRequest;
import org.apache.zookeeper.proto.SetDataRequest;
import org.apache.zookeeper.server.ByteBufferInputStream;
import org.apache.zookeeper.server.DataTree.ProcessTxnResult;
import org.apache.zookeeper.server.Request;
+import org.apache.zookeeper.server.auth.AuthenticationProvider;
+import org.apache.zookeeper.server.auth.DigestAuthenticationProvider;
+import org.apache.zookeeper.server.auth.IPAuthenticationProvider;
+import org.apache.zookeeper.server.auth.ProviderRegistry;
+import org.apache.zookeeper.server.auth.SASLAuthenticationProvider;
+import org.apache.zookeeper.server.auth.X509AuthenticationProvider;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -43,6 +55,7 @@
*/
public final class AuditHelper {
private static final Logger LOG = LoggerFactory.getLogger(AuditHelper.class);
+ private static final String REDACTED = "[redacted]";
public static void addAuditLog(Request request, ProcessTxnResult rc) {
addAuditLog(request, rc, false);
@@ -53,174 +66,334 @@ public static void addAuditLog(Request request, ProcessTxnResult rc) {
*
* @param request user request
* @param txnResult ProcessTxnResult
- * @param failedTxn whether audit is being done failed transaction for normal transaction
+ * @param failedTxn whether the transaction was rejected before applying the requested operation
*/
public static void addAuditLog(Request request, ProcessTxnResult txnResult, boolean failedTxn) {
if (!ZKAuditProvider.isAuditEnabled()) {
return;
}
- String op = null;
- //For failed transaction rc.path is null
- String path = txnResult.path;
- String acls = null;
- String createMode = null;
try {
- switch (request.type) {
- case ZooDefs.OpCode.create:
- case ZooDefs.OpCode.create2:
- case ZooDefs.OpCode.createContainer:
- op = AuditConstants.OP_CREATE;
- if (failedTxn) {
- CreateRequest createRequest = new CreateRequest();
- deserialize(request, createRequest);
- path = createRequest.getPath();
- createMode =
- getCreateMode(createRequest);
- } else {
- createMode = getCreateMode(request);
- }
- break;
- case ZooDefs.OpCode.delete:
- case ZooDefs.OpCode.deleteContainer:
- op = AuditConstants.OP_DELETE;
- if (failedTxn) {
- DeleteRequest deleteRequest = new DeleteRequest();
- deserialize(request, deleteRequest);
- path = deleteRequest.getPath();
- }
- break;
- case ZooDefs.OpCode.setData:
- op = AuditConstants.OP_SETDATA;
- if (failedTxn) {
- SetDataRequest setDataRequest = new SetDataRequest();
- deserialize(request, setDataRequest);
- path = setDataRequest.getPath();
- }
- break;
- case ZooDefs.OpCode.setACL:
- op = AuditConstants.OP_SETACL;
- if (failedTxn) {
- SetACLRequest setACLRequest = new SetACLRequest();
- deserialize(request, setACLRequest);
- path = setACLRequest.getPath();
- acls = ZKUtil.aclToString(setACLRequest.getAcl());
- } else {
- acls = getACLs(request);
+ if (request.type == OpCode.multi) {
+ logMultiOperation(request, txnResult, failedTxn);
+ return;
+ }
+ String operation = operationFor(request.type);
+ if (operation == null) {
+ return;
+ }
+ Integer error = errorCode(request, txnResult, failedTxn);
+ Outcome outcome = outcome(txnResult, failedTxn, error);
+ RequestMetadata metadata = new RequestMetadata();
+ boolean enhanced = ZKAuditProvider.isEnhancedAuditEnabled();
+ if (isCreate(request.type) || request.type == OpCode.setACL
+ || (request.type != OpCode.reconfig && (enhanced || outcome != Outcome.COMMITTED))) {
+ try {
+ metadata = metadata(request.type, readRequestRecord(request), enhanced);
+ } catch (IOException e) {
+ auditError(request.type, e);
+ }
+ }
+ String path = txnResult == null ? null : txnResult.path;
+ if (outcome != Outcome.COMMITTED || path == null) {
+ path = metadata.path == null ? path : metadata.path;
+ }
+ log(request, txnResult, path, operation, metadata, result(outcome), error, outcome, null, enhanced);
+ } catch (RuntimeException e) {
+ auditError(request.type, e);
+ }
+ }
+
+ private static Record readRequestRecord(Request request) throws IOException {
+ Record record;
+ switch (request.type) {
+ case OpCode.create:
+ case OpCode.create2:
+ case OpCode.createContainer:
+ record = new CreateRequest();
+ break;
+ case OpCode.createTTL:
+ record = new CreateTTLRequest();
+ break;
+ case OpCode.delete:
+ case OpCode.deleteContainer:
+ record = new DeleteRequest();
+ break;
+ case OpCode.setData:
+ record = new SetDataRequest();
+ break;
+ case OpCode.setACL:
+ record = new SetACLRequest();
+ break;
+ case OpCode.multi:
+ record = new MultiOperationRecord();
+ break;
+ default:
+ throw new IOException("Unsupported audit request type: " + request.type);
+ }
+ if (request.request == null) {
+ throw new IOException("Audit request record is unavailable");
+ }
+ ByteBuffer buffer = request.request.duplicate();
+ buffer.rewind();
+ ByteBufferInputStream.byteBuffer2Record(buffer, record);
+ return record;
+ }
+
+ private static void logMultiOperation(Request request, ProcessTxnResult rc, boolean failedTxn) {
+ boolean enhanced = ZKAuditProvider.isEnhancedAuditEnabled();
+ Integer error = errorCode(request, rc, failedTxn);
+ boolean failed = failedTxn || (error != null && error != Code.OK.intValue());
+ int failureIndex = -1;
+ if (rc != null && rc.multiResult != null) {
+ for (int index = 0; index < rc.multiResult.size(); index++) {
+ ProcessTxnResult subResult = rc.multiResult.get(index);
+ if (subResult != null && (subResult.type == OpCode.error || subResult.err != Code.OK.intValue())) {
+ failed = true;
+ if (failureIndex < 0 && subResult.err != Code.OK.intValue()) {
+ failureIndex = index;
}
- break;
- case ZooDefs.OpCode.multi:
- if (failedTxn) {
- op = AuditConstants.OP_MULTI_OP;
- } else {
- logMultiOperation(request, txnResult);
- //operation si already logged
- return;
+ if (error == null || error == Code.OK.intValue()) {
+ error = subResult.err;
}
+ }
+ }
+ }
+ // Emit the failed parent before decoding so malformed metadata cannot hide the failure.
+ if (failed) {
+ log(request, rc, rc == null ? null : rc.path, AuditConstants.OP_MULTI_OP,
+ new RequestMetadata(), Result.FAILURE, error, Outcome.FAILED, null, enhanced);
+ if (!enhanced) {
+ return;
+ }
+ }
+ MultiOperationRecord multiRequest;
+ try {
+ multiRequest = (MultiOperationRecord) readRequestRecord(request);
+ } catch (IOException e) {
+ auditError(request.type, e);
+ return;
+ }
+ boolean complete = rc != null && rc.type == OpCode.multi && rc.multiResult != null
+ && rc.multiResult.size() == multiRequest.size();
+ if (complete) {
+ int index = 0;
+ for (Op op : multiRequest) {
+ ProcessTxnResult subResult = rc.multiResult.get(index++);
+ if (subResult == null || (subResult.type != OpCode.error
+ && subResult.type != op.getType()
+ && !(subResult.type == OpCode.create2 && op.getType() == OpCode.create))) {
+ complete = false;
break;
- case ZooDefs.OpCode.reconfig:
- op = AuditConstants.OP_RECONFIG;
- break;
- default:
- //Not an audit log operation
- return;
+ }
+ }
+ }
+ if (!complete) {
+ auditError(request.type, new IOException("Audit multi request/result positions do not match"));
+ }
+ int index = 0;
+ for (Op op : multiRequest) {
+ ProcessTxnResult subResult = complete ? rc.multiResult.get(index) : null;
+ String operation = operationFor(op.getType());
+ if (operation != null) {
+ Outcome outcome;
+ if (!complete) {
+ outcome = Outcome.UNKNOWN;
+ } else if (failed) {
+ // A zero-code error member was rolled back, not committed. Later members are not applied.
+ outcome = subResult.err == Code.OK.intValue()
+ || (index > failureIndex && subResult.err == Code.RUNTIMEINCONSISTENCY.intValue())
+ ? Outcome.ROLLED_BACK : Outcome.FAILED;
+ } else {
+ outcome = Outcome.COMMITTED;
+ }
+ RequestMetadata metadata = metadata(op.getType(), op.toRequestRecord(), enhanced);
+ String path = outcome == Outcome.COMMITTED ? subResult.path : op.getPath();
+ log(request, rc, path, operation, metadata, failed ? Result.FAILURE : result(outcome),
+ subResult == null ? null : subResult.err, outcome, index, enhanced);
}
- Result result = getResult(txnResult, failedTxn);
- log(request, path, op, acls, createMode, result);
- } catch (Throwable e) {
- LOG.error("Failed to audit log request {}", request.type, e);
+ index++;
}
}
- private static void deserialize(Request request, Record record) throws IOException {
- request.request.rewind();
- ByteBufferInputStream.byteBuffer2Record(request.request.slice(), record);
+ private static RequestMetadata metadata(int type, Record record, boolean enhanced) {
+ RequestMetadata metadata = new RequestMetadata();
+ switch (type) {
+ case OpCode.create:
+ case OpCode.create2:
+ case OpCode.createContainer:
+ CreateRequest create = (CreateRequest) record;
+ metadata.path = create.getPath();
+ metadata.dataLength = length(create.getData());
+ metadata.createMode = createMode(type, create.getFlags());
+ break;
+ case OpCode.createTTL:
+ CreateTTLRequest ttl = (CreateTTLRequest) record;
+ metadata.path = ttl.getPath();
+ metadata.dataLength = length(ttl.getData());
+ metadata.createMode = createMode(type, ttl.getFlags());
+ break;
+ case OpCode.setData:
+ SetDataRequest setData = (SetDataRequest) record;
+ metadata.path = setData.getPath();
+ metadata.dataLength = length(setData.getData());
+ break;
+ case OpCode.delete:
+ case OpCode.deleteContainer:
+ metadata.path = ((DeleteRequest) record).getPath();
+ break;
+ case OpCode.setACL:
+ SetACLRequest setAcl = (SetACLRequest) record;
+ metadata.path = setAcl.getPath();
+ if (setAcl.getAcl() != null) {
+ metadata.acl = enhanced ? safeAclToString(setAcl.getAcl()) : ZKUtil.aclToString(setAcl.getAcl());
+ }
+ break;
+ default:
+ break;
+ }
+ return metadata;
}
- private static Result getResult(ProcessTxnResult rc, boolean failedTxn) {
- if (failedTxn) {
- return Result.FAILURE;
- } else {
- return rc.err == KeeperException.Code.OK.intValue() ? Result.SUCCESS : Result.FAILURE;
+ private static String safeAclToString(List acls) {
+ StringBuilder value = new StringBuilder();
+ for (ACL acl : acls) {
+ Id id = acl.getId();
+ String user = REDACTED;
+ if ("world".equals(id.getScheme()) && "anyone".equals(id.getId())) {
+ user = "anyone";
+ } else if ("digest".equals(id.getScheme())) {
+ AuthenticationProvider provider = ProviderRegistry.getProvider(id.getScheme());
+ if (provider != null && provider.getClass() == DigestAuthenticationProvider.class
+ && id.getId() != null && provider.isValid(id.getId())) {
+ user = provider.getUserName(id.getId());
+ }
+ }
+ value.append(id.getScheme()).append(':').append(user).append(':')
+ .append(ZKUtil.getPermString(acl.getPerms()));
}
+ return value.toString();
}
- private static void logMultiOperation(Request request, ProcessTxnResult rc) throws IOException, KeeperException {
- Map createModes = AuditHelper.getCreateModes(request);
- boolean multiFailed = false;
- for (ProcessTxnResult subTxnResult : rc.multiResult) {
- switch (subTxnResult.type) {
- case ZooDefs.OpCode.create:
- case ZooDefs.OpCode.create2:
- case ZooDefs.OpCode.createTTL:
- case ZooDefs.OpCode.createContainer:
- log(request, subTxnResult.path, AuditConstants.OP_CREATE, null,
- createModes.get(subTxnResult.path), Result.SUCCESS);
- break;
- case ZooDefs.OpCode.delete:
- case ZooDefs.OpCode.deleteContainer:
- log(request, subTxnResult.path, AuditConstants.OP_DELETE, null,
- null, Result.SUCCESS);
- break;
- case ZooDefs.OpCode.setData:
- log(request, subTxnResult.path, AuditConstants.OP_SETDATA, null,
- null, Result.SUCCESS);
- break;
- case ZooDefs.OpCode.error:
- multiFailed = true;
- break;
- default:
- // Do nothing, it ok, we do not log all multi operations
+ private static String getUsers(Request request, boolean enhanced) {
+ if (!enhanced) {
+ return request.getUsers();
+ }
+ if (request.authInfo == null) {
+ return null;
+ }
+ StringBuilder users = new StringBuilder();
+ boolean first = true;
+ for (Id id : request.authInfo) {
+ if (!first) {
+ users.append(',');
}
+ first = false;
+ users.append(safeUser(id));
+ }
+ return users.toString();
+ }
+
+ private static String safeUser(Id id) {
+ if (id == null || id.getScheme() == null || id.getId() == null) {
+ return REDACTED;
+ }
+ AuthenticationProvider provider = ProviderRegistry.getProvider(id.getScheme());
+ if (provider == null) {
+ return REDACTED;
}
- if (multiFailed) {
- log(request, rc.path, AuditConstants.OP_MULTI_OP, null,
- null, Result.FAILURE);
+ Class> providerClass = provider.getClass();
+ // Custom providers, including subclasses, do not establish safe identity representations.
+ if ((providerClass == DigestAuthenticationProvider.class
+ || providerClass == IPAuthenticationProvider.class
+ || providerClass == SASLAuthenticationProvider.class
+ || providerClass == X509AuthenticationProvider.class)
+ && provider.isValid(id.getId())) {
+ return provider.getUserName(id.getId());
}
+ return REDACTED;
}
- private static void log(Request request, String path, String op, String acls, String createMode, Result result) {
- log(request.getUsers(), op, path, acls, createMode,
- request.cnxn.getSessionIdHex(), request.cnxn.getHostAddress(), result);
+ private static String createMode(int type, int flags) {
+ try {
+ return CreateMode.fromFlag(flags).name().toLowerCase(Locale.ROOT);
+ } catch (KeeperException e) {
+ auditError(type, e);
+ return null;
+ }
}
- private static void log(String user, String operation, String znode, String acl,
- String createMode, String session, String ip, Result result) {
- ZKAuditProvider.log(user, operation, znode, acl, createMode, session, ip, result);
+ private static int length(byte[] data) {
+ return data == null ? 0 : data.length;
}
- private static String getACLs(Request request) throws IOException {
- SetACLRequest setACLRequest = new SetACLRequest();
- deserialize(request, setACLRequest);
- return ZKUtil.aclToString(setACLRequest.getAcl());
+ private static Integer errorCode(Request request, ProcessTxnResult rc, boolean failedTxn) {
+ if (failedTxn && request.getException() != null) {
+ return request.getException().code().intValue();
+ }
+ return rc == null || rc.type == 0 ? null : rc.err;
}
- private static String getCreateMode(Request request) throws IOException, KeeperException {
- CreateRequest createRequest = new CreateRequest();
- deserialize(request, createRequest);
- return getCreateMode(createRequest);
+ private static Outcome outcome(ProcessTxnResult rc, boolean failedTxn, Integer error) {
+ if (failedTxn || (rc != null && rc.type == OpCode.error)
+ || (error != null && error != Code.OK.intValue())) {
+ return Outcome.FAILED;
+ }
+ return error == null ? Outcome.UNKNOWN : Outcome.COMMITTED;
}
- private static String getCreateMode(CreateRequest createRequest) throws KeeperException {
- return CreateMode.fromFlag(createRequest.getFlags()).toString().toLowerCase();
+ private static Result result(Outcome outcome) {
+ return outcome == Outcome.COMMITTED ? Result.SUCCESS : outcome == Outcome.UNKNOWN ? Result.INVOKED : Result.FAILURE;
}
- private static Map getCreateModes(Request request)
- throws IOException, KeeperException {
- Map createModes = new HashMap<>();
- if (!ZKAuditProvider.isAuditEnabled()) {
- return createModes;
+ private static String operationFor(int type) {
+ switch (type) {
+ case OpCode.create:
+ case OpCode.create2:
+ case OpCode.createTTL:
+ case OpCode.createContainer:
+ return AuditConstants.OP_CREATE;
+ case OpCode.delete:
+ case OpCode.deleteContainer:
+ return AuditConstants.OP_DELETE;
+ case OpCode.setData:
+ return AuditConstants.OP_SETDATA;
+ case OpCode.setACL:
+ return AuditConstants.OP_SETACL;
+ case OpCode.reconfig:
+ return AuditConstants.OP_RECONFIG;
+ default:
+ return null;
}
- MultiOperationRecord multiRequest = new MultiOperationRecord();
- deserialize(request, multiRequest);
- for (Op op : multiRequest) {
- if (op.getType() == ZooDefs.OpCode.create || op.getType() == ZooDefs.OpCode.create2
- || op.getType() == ZooDefs.OpCode.createContainer) {
- CreateRequest requestRecord = (CreateRequest) op.toRequestRecord();
- createModes.put(requestRecord.getPath(),
- getCreateMode(requestRecord));
- }
+ }
+
+ private static boolean isCreate(int type) {
+ return type == OpCode.create || type == OpCode.create2 || type == OpCode.createTTL || type == OpCode.createContainer;
+ }
+
+ private static void log(Request request, ProcessTxnResult rc, String path, String operation,
+ RequestMetadata metadata, Result result, Integer error, Outcome outcome,
+ Integer index, boolean enhanced) {
+ Long zxid = null;
+ if (request.getHdr() != null) {
+ zxid = request.getHdr().getZxid();
+ } else if (rc != null && rc.type != 0) {
+ zxid = rc.zxid;
+ } else if (request.zxid >= 0) {
+ zxid = request.zxid;
}
- return createModes;
+ ZKAuditProvider.log(getUsers(request, enhanced), operation, path, metadata.acl, metadata.createMode,
+ request.cnxn.getSessionIdHex(), request.cnxn.getHostAddress(), result,
+ metadata.dataLength, error, outcome, request.cxid, zxid, index, enhanced);
}
+ private static void auditError(int type, Exception e) {
+ ZKAuditProvider.reportAuditError(LOG, "Failed to audit log request {}", type, e);
+ }
+
+ private static final class RequestMetadata {
+ private String path;
+ private String acl;
+ private String createMode;
+ private Integer dataLength;
+ }
}
diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/audit/ZKAuditProvider.java b/zookeeper-server/src/main/java/org/apache/zookeeper/audit/ZKAuditProvider.java
index 864d2366c97..f05062f9d1c 100644
--- a/zookeeper-server/src/main/java/org/apache/zookeeper/audit/ZKAuditProvider.java
+++ b/zookeeper-server/src/main/java/org/apache/zookeeper/audit/ZKAuditProvider.java
@@ -19,13 +19,17 @@
import static org.apache.zookeeper.audit.AuditEvent.FieldName;
import java.lang.reflect.Constructor;
+import java.util.Locale;
+import org.apache.zookeeper.audit.AuditEvent.Outcome;
import org.apache.zookeeper.audit.AuditEvent.Result;
import org.apache.zookeeper.server.ServerCnxnFactory;
+import org.apache.zookeeper.server.ServerMetrics;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
public class ZKAuditProvider {
static final String AUDIT_ENABLE = "zookeeper.audit.enable";
+ static final String AUDIT_ENHANCED_ENABLE = "zookeeper.audit.enhanced.enable";
static final String AUDIT_IMPL_CLASS = "zookeeper.audit.impl.class";
private static final Logger LOG = LoggerFactory.getLogger(ZKAuditProvider.class);
// By default audit logging is disabled
@@ -66,9 +70,39 @@ public static boolean isAuditEnabled() {
return auditEnabled;
}
+ public static boolean isEnhancedAuditEnabled() {
+ return auditEnabled && Boolean.getBoolean(AUDIT_ENHANCED_ENABLE);
+ }
+
public static void log(String user, String operation, String znode, String acl,
String createMode, String session, String ip, Result result) {
- auditLogger.logAuditEvent(createLogEvent(user, operation, znode, acl, createMode, session, ip, result));
+ log(user, operation, znode, acl, createMode, session, ip, result, null, null, null, null, null, null);
+ }
+
+ /**
+ * Logs optional operation metadata only when enhanced audit logging is enabled.
+ * Null metadata is unavailable, not zero. Callers must sanitize ACL identities.
+ */
+ public static void log(String user, String operation, String znode, String acl,
+ String createMode, String session, String ip, Result result,
+ Integer dataLength, Integer errorCode, Outcome outcome,
+ Integer cxid, Long zxid, Integer multiIndex) {
+ if (!isAuditEnabled()) {
+ return;
+ }
+ log(user, operation, znode, acl, createMode, session, ip, result,
+ dataLength, errorCode, outcome, cxid, zxid, multiIndex, isEnhancedAuditEnabled());
+ }
+
+ static void log(String user, String operation, String znode, String acl,
+ String createMode, String session, String ip, Result result,
+ Integer dataLength, Integer errorCode, Outcome outcome,
+ Integer cxid, Long zxid, Integer multiIndex, boolean enhanced) {
+ if (!isAuditEnabled()) {
+ return;
+ }
+ logAuditEvent(createLogEvent(user, operation, znode, acl, createMode, session, ip, result,
+ dataLength, errorCode, outcome, cxid, zxid, multiIndex, enhanced));
}
/**
@@ -78,6 +112,7 @@ static AuditEvent createLogEvent(String user, String operation, Result result) {
AuditEvent event = new AuditEvent(result);
event.addEntry(FieldName.USER, user);
event.addEntry(FieldName.OPERATION, operation);
+ addMetadata(event, null, null, null, null, null, null, isEnhancedAuditEnabled());
return event;
}
@@ -86,6 +121,22 @@ static AuditEvent createLogEvent(String user, String operation, Result result) {
*/
static AuditEvent createLogEvent(String user, String operation, String znode, String acl,
String createMode, String session, String ip, Result result) {
+ return createLogEvent(user, operation, znode, acl, createMode, session, ip, result,
+ null, null, null, null, null, null);
+ }
+
+ static AuditEvent createLogEvent(String user, String operation, String znode, String acl,
+ String createMode, String session, String ip, Result result,
+ Integer dataLength, Integer errorCode, Outcome outcome,
+ Integer cxid, Long zxid, Integer multiIndex) {
+ return createLogEvent(user, operation, znode, acl, createMode, session, ip, result,
+ dataLength, errorCode, outcome, cxid, zxid, multiIndex, isEnhancedAuditEnabled());
+ }
+
+ private static AuditEvent createLogEvent(String user, String operation, String znode, String acl,
+ String createMode, String session, String ip, Result result,
+ Integer dataLength, Integer errorCode, Outcome outcome,
+ Integer cxid, Long zxid, Integer multiIndex, boolean enhanced) {
AuditEvent event = new AuditEvent(result);
event.addEntry(FieldName.SESSION, session);
event.addEntry(FieldName.USER, user);
@@ -94,9 +145,48 @@ static AuditEvent createLogEvent(String user, String operation, String znode, St
event.addEntry(FieldName.ZNODE, znode);
event.addEntry(FieldName.ZNODE_TYPE, createMode);
event.addEntry(FieldName.ACL, acl);
+ addMetadata(event, dataLength, errorCode, outcome, cxid, zxid, multiIndex, enhanced);
return event;
}
+ private static void addMetadata(AuditEvent event, Integer dataLength, Integer errorCode, Outcome outcome,
+ Integer cxid, Long zxid, Integer multiIndex, boolean enhanced) {
+ if (enhanced) {
+ event.addEntry(FieldName.SCHEMA_VERSION, AuditConstants.SCHEMA_VERSION);
+ event.addEntry(FieldName.DATA_LENGTH, valueOf(dataLength));
+ event.addEntry(FieldName.ERROR_CODE, valueOf(errorCode));
+ event.addEntry(FieldName.OUTCOME, (outcome == null ? Outcome.UNKNOWN : outcome).name().toLowerCase(Locale.ROOT));
+ event.addEntry(FieldName.CXID, valueOf(cxid));
+ event.addEntry(FieldName.ZXID, valueOf(zxid));
+ event.addEntry(FieldName.MULTI_INDEX, valueOf(multiIndex));
+ }
+ }
+
+ private static String valueOf(Number value) {
+ return value == null ? null : value.toString();
+ }
+
+ private static void logAuditEvent(AuditEvent event) {
+ try {
+ auditLogger.logAuditEvent(event);
+ } catch (RuntimeException e) {
+ reportAuditError(LOG, "Failed to write audit log for operation {}", event.getValue(FieldName.OPERATION), e);
+ }
+ }
+
+ static void reportAuditError(Logger logger, String message, Object context, Exception error) {
+ try {
+ ServerMetrics.getMetrics().AUDIT_ERRORS.add(1);
+ } catch (RuntimeException ignored) {
+ // A failed metrics backend must not prevent the diagnostic or change the operation's result.
+ }
+ try {
+ logger.error(message, context, error);
+ } catch (RuntimeException ignored) {
+ // Reporting is best-effort; retrying through the same failing logger could escape or recurse.
+ }
+ }
+
/**
* Add audit log for server start and register server stop log.
*/
@@ -119,7 +209,7 @@ public static void addServerStartFailureAuditLog() {
}
private static void log(String user, String operation, Result result) {
- auditLogger.logAuditEvent(createLogEvent(user, operation, result));
+ logAuditEvent(createLogEvent(user, operation, result));
}
/**
diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/DataTree.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/DataTree.java
index 1a5d1304120..80d0173afe6 100644
--- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/DataTree.java
+++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/DataTree.java
@@ -59,6 +59,7 @@
import org.apache.zookeeper.ZooDefs;
import org.apache.zookeeper.ZooDefs.OpCode;
import org.apache.zookeeper.audit.AuditConstants;
+import org.apache.zookeeper.audit.AuditEvent.Outcome;
import org.apache.zookeeper.audit.AuditEvent.Result;
import org.apache.zookeeper.audit.ZKAuditProvider;
import org.apache.zookeeper.common.PathTrie;
@@ -1451,15 +1452,11 @@ void deleteNodes(long session, long zxid, Iterable paths2Delete) {
path, sessionHex);
}
if (ZKAuditProvider.isAuditEnabled()) {
- if (deleted) {
- ZKAuditProvider.log(ZKAuditProvider.getZKUser(),
- AuditConstants.OP_DEL_EZNODE_EXP, path, null, null,
- sessionHex, null, Result.SUCCESS);
- } else {
- ZKAuditProvider.log(ZKAuditProvider.getZKUser(),
- AuditConstants.OP_DEL_EZNODE_EXP, path, null, null,
- sessionHex, null, Result.FAILURE);
- }
+ ZKAuditProvider.log(ZKAuditProvider.getZKUser(),
+ AuditConstants.OP_DEL_EZNODE_EXP, path, null, null,
+ sessionHex, null, deleted ? Result.SUCCESS : Result.FAILURE,
+ null, deleted ? Code.OK.intValue() : Code.NONODE.intValue(),
+ deleted ? Outcome.COMMITTED : Outcome.FAILED, null, zxid, null);
}
}
}
diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/ServerMetrics.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/ServerMetrics.java
index c9cfcc12a52..4b5ba818106 100644
--- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/ServerMetrics.java
+++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/ServerMetrics.java
@@ -227,6 +227,7 @@ private ServerMetrics(MetricsProvider metricsProvider) {
STALE_REQUESTS = metricsContext.getCounter("stale_requests");
STALE_REQUESTS_DROPPED = metricsContext.getCounter("stale_requests_dropped");
STALE_REPLIES = metricsContext.getCounter("stale_replies");
+ AUDIT_ERRORS = metricsContext.getCounter("audit_errors");
REQUEST_THROTTLE_WAIT_COUNT = metricsContext.getCounter("request_throttle_wait_count");
LARGE_REQUESTS_REJECTED = metricsContext.getCounter("large_requests_rejected");
@@ -443,6 +444,7 @@ private ServerMetrics(MetricsProvider metricsProvider) {
public final Counter STALE_REQUESTS;
public final Counter STALE_REQUESTS_DROPPED;
public final Counter STALE_REPLIES;
+ public final Counter AUDIT_ERRORS;
public final Counter REQUEST_THROTTLE_WAIT_COUNT;
public final Counter LARGE_REQUESTS_REJECTED;
diff --git a/zookeeper-server/src/test/java/org/apache/zookeeper/audit/AuditEventTest.java b/zookeeper-server/src/test/java/org/apache/zookeeper/audit/AuditEventTest.java
index 02d9ac0bb85..d5f371a7a6a 100644
--- a/zookeeper-server/src/test/java/org/apache/zookeeper/audit/AuditEventTest.java
+++ b/zookeeper-server/src/test/java/org/apache/zookeeper/audit/AuditEventTest.java
@@ -18,6 +18,8 @@
package org.apache.zookeeper.audit;
import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertFalse;
+import java.util.Locale;
import org.apache.zookeeper.audit.AuditEvent.Result;
import org.junit.Test;
@@ -42,4 +44,55 @@ public void testFormatShouldIgnoreKeyIfValueIsNull() {
String expected = "operation=Value2\tresult=success";
assertEquals(expected, actual);
}
+
+ @Test
+ public void testProtocolNamesDoNotDependOnDefaultLocale() {
+ Locale previous = Locale.getDefault();
+ try {
+ Locale.setDefault(new Locale("tr", "TR"));
+ AuditEvent event = new AuditEvent(Result.FAILURE);
+ event.addEntry(AuditEvent.FieldName.IP, "127.0.0.1");
+ assertEquals("ip=127.0.0.1\tresult=failure", event.toString());
+ assertEquals("result=invoked", new AuditEvent(Result.INVOKED).toString());
+ } finally {
+ Locale.setDefault(previous);
+ }
+ }
+
+ @Test
+ public void testEnhancedFormattingEscapesValuesReversibly() {
+ String previousAudit = System.getProperty(ZKAuditProvider.AUDIT_ENABLE);
+ String previousEnhanced = System.getProperty(AuditHelperTest.ENHANCED_ENABLE);
+ try {
+ System.setProperty(ZKAuditProvider.AUDIT_ENABLE, "true");
+ System.setProperty(AuditHelperTest.ENHANCED_ENABLE, "true");
+ AuditEvent event = ZKAuditProvider.createLogEvent("team\tname\r\n\\t", "setData",
+ "/name=value\\child", null, null, null, null, Result.FAILURE);
+ String log = event.toString();
+ assertEquals("2", AuditHelperTest.fields(log).get("schema_version"));
+ assertEquals("team\\tname\\r\\n\\\\t", AuditHelperTest.fields(log).get("user"));
+ assertEquals("/name=value\\\\child", AuditHelperTest.fields(log).get("znode"));
+ assertEquals("team\tname\r\n\\t", event.getValue(AuditEvent.FieldName.USER));
+ assertFalse(log.contains("\n"));
+ assertFalse(log.contains("\r"));
+ } finally {
+ AuditHelperTest.restoreProperty(ZKAuditProvider.AUDIT_ENABLE, previousAudit);
+ AuditHelperTest.restoreProperty(AuditHelperTest.ENHANCED_ENABLE, previousEnhanced);
+ }
+ }
+
+ @Test
+ public void testLegacyValuesAreNotEscaped() {
+ AuditEvent event = new AuditEvent(Result.SUCCESS);
+ event.addEntry(AuditEvent.FieldName.USER, "team\tname\n\\t");
+ assertEquals("user=team\tname\n\\t\tresult=success", event.toString());
+ }
+
+ @Test
+ public void testSchemaMarkerEscapesPreviouslyAddedValues() {
+ AuditEvent event = new AuditEvent(Result.SUCCESS);
+ event.addEntry(AuditEvent.FieldName.USER, "team\tname\n");
+ event.addEntry(AuditEvent.FieldName.SCHEMA_VERSION, "2");
+ assertEquals("user=team\\tname\\n\tschema_version=2\tresult=success", event.toString());
+ }
}
diff --git a/zookeeper-server/src/test/java/org/apache/zookeeper/audit/AuditHelperTest.java b/zookeeper-server/src/test/java/org/apache/zookeeper/audit/AuditHelperTest.java
new file mode 100644
index 00000000000..2dc2ab9f47a
--- /dev/null
+++ b/zookeeper-server/src/test/java/org/apache/zookeeper/audit/AuditHelperTest.java
@@ -0,0 +1,855 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you 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
+ *
+ * http://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 org.apache.zookeeper.audit;
+
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertNotNull;
+import static org.junit.Assert.assertNull;
+import static org.junit.Assert.assertTrue;
+import static org.junit.Assert.fail;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
+import ch.qos.logback.classic.Level;
+import ch.qos.logback.classic.Logger;
+import ch.qos.logback.classic.spi.ILoggingEvent;
+import ch.qos.logback.core.AppenderBase;
+import java.io.ByteArrayOutputStream;
+import java.io.IOException;
+import java.lang.reflect.Field;
+import java.nio.ByteBuffer;
+import java.nio.charset.StandardCharsets;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.LinkedHashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+import org.apache.jute.BinaryOutputArchive;
+import org.apache.jute.Record;
+import org.apache.zookeeper.CreateMode;
+import org.apache.zookeeper.KeeperException;
+import org.apache.zookeeper.KeeperException.Code;
+import org.apache.zookeeper.MultiOperationRecord;
+import org.apache.zookeeper.Op;
+import org.apache.zookeeper.ZooDefs;
+import org.apache.zookeeper.ZooDefs.OpCode;
+import org.apache.zookeeper.audit.AuditEvent.Result;
+import org.apache.zookeeper.data.ACL;
+import org.apache.zookeeper.data.Id;
+import org.apache.zookeeper.metrics.Counter;
+import org.apache.zookeeper.metrics.MetricsUtils;
+import org.apache.zookeeper.proto.CreateRequest;
+import org.apache.zookeeper.proto.CreateTTLRequest;
+import org.apache.zookeeper.proto.SetACLRequest;
+import org.apache.zookeeper.proto.SetDataRequest;
+import org.apache.zookeeper.server.DataTree;
+import org.apache.zookeeper.server.DataTree.ProcessTxnResult;
+import org.apache.zookeeper.server.Request;
+import org.apache.zookeeper.server.ServerCnxn;
+import org.apache.zookeeper.server.ServerMetrics;
+import org.apache.zookeeper.server.auth.AuthenticationProvider;
+import org.apache.zookeeper.server.auth.DigestAuthenticationProvider;
+import org.apache.zookeeper.server.auth.ProviderRegistry;
+import org.apache.zookeeper.txn.CheckVersionTxn;
+import org.apache.zookeeper.txn.CloseSessionTxn;
+import org.apache.zookeeper.txn.CreateTTLTxn;
+import org.apache.zookeeper.txn.CreateTxn;
+import org.apache.zookeeper.txn.DeleteTxn;
+import org.apache.zookeeper.txn.ErrorTxn;
+import org.apache.zookeeper.txn.MultiTxn;
+import org.apache.zookeeper.txn.SetACLTxn;
+import org.apache.zookeeper.txn.SetDataTxn;
+import org.apache.zookeeper.txn.Txn;
+import org.apache.zookeeper.txn.TxnHeader;
+import org.junit.After;
+import org.junit.Before;
+import org.junit.Test;
+import org.slf4j.LoggerFactory;
+
+public class AuditHelperTest {
+ static final String ENHANCED_ENABLE = "zookeeper.audit.enhanced.enable";
+ private static final long SESSION = 0x123;
+ private AuditCapture capture;
+ private DataTree tree;
+ private ServerCnxn cnxn;
+ private String previousEnhanced;
+ private String previousAudit;
+ private String previousExtendedTypes;
+
+ @Before
+ public void setUp() {
+ previousAudit = System.getProperty(ZKAuditProvider.AUDIT_ENABLE);
+ previousEnhanced = System.getProperty(ENHANCED_ENABLE);
+ previousExtendedTypes = System.getProperty("zookeeper.extendedTypesEnabled");
+ System.setProperty(ZKAuditProvider.AUDIT_ENABLE, "true");
+ System.setProperty(ENHANCED_ENABLE, "true");
+ System.setProperty("zookeeper.extendedTypesEnabled", "true");
+ assertTrue(ZKAuditProvider.isAuditEnabled());
+ capture = new AuditCapture();
+ tree = new DataTree();
+ cnxn = mock(ServerCnxn.class);
+ when(cnxn.getSessionIdHex()).thenReturn("0x123");
+ when(cnxn.getHostAddress()).thenReturn("127.0.0.1");
+ }
+
+ @After
+ public void tearDown() {
+ capture.close();
+ restoreProperty(ZKAuditProvider.AUDIT_ENABLE, previousAudit);
+ restoreProperty(ENHANCED_ENABLE, previousEnhanced);
+ restoreProperty("zookeeper.extendedTypesEnabled", previousExtendedTypes);
+ }
+
+ @Test
+ public void testCreateAndSetDataLengthsArePayloadBytes() throws Exception {
+ byte[][] data = {null, new byte[0], new byte[1], new byte[1024],
+ new byte[65536], new byte[1048575], "\u00e9\ud83d\ude00".getBytes(StandardCharsets.UTF_8)};
+ String[] lengths = {"0", "0", "1", "1024", "65536", "1048575", "6"};
+ for (int i = 0; i < data.length; i++) {
+ String path = "/length-" + i;
+ Request create = request(OpCode.create, createRecord(path, data[i], CreateMode.PERSISTENT));
+ ProcessTxnResult created = apply(create, OpCode.create, createTxn(path, data[i], false));
+ assertEquals(0, created.err);
+ AuditHelper.addAuditLog(create, created);
+ Map fields = fields(capture.read(1).get(0));
+ assertWrite(fields, "create", path, lengths[i], "committed", "0");
+ assertEquals("41", fields.get("cxid"));
+ assertEquals("66", fields.get("zxid"));
+ assertNull(fields.get("multi_index"));
+
+ Request setData = request(OpCode.setData, new SetDataRequest(path, data[i], -1));
+ ProcessTxnResult changed = apply(setData, OpCode.setData, new SetDataTxn(path, data[i], 1));
+ assertEquals(0, changed.err);
+ AuditHelper.addAuditLog(setData, changed);
+ assertWrite(fields(capture.read(1).get(0)), "setData", path, lengths[i], "committed", "0");
+ }
+ }
+
+ @Test
+ public void testCreateTtlUsesTypedRecord() throws Exception {
+ byte[] data = {1, 2, 3};
+ Request request = request(OpCode.createTTL, new CreateTTLRequest(
+ "/ttl", data, ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT_WITH_TTL.toFlag(), 60000));
+ ProcessTxnResult result = apply(request, OpCode.createTTL,
+ new CreateTTLTxn("/ttl", data, ZooDefs.Ids.OPEN_ACL_UNSAFE, -1, 60000));
+ assertEquals(0, result.err);
+ AuditHelper.addAuditLog(request, result);
+ Map fields = fields(capture.read(1).get(0));
+ assertWrite(fields, "create", "/ttl", "3", "committed", "0");
+ assertEquals("persistent_with_ttl", fields.get("znode_type"));
+ }
+
+ @Test
+ public void testFailureUsesAttemptedLengthAndReplyException() throws Exception {
+ Request request = request(OpCode.create, createRecord("/denied", new byte[7], CreateMode.PERSISTENT));
+ ProcessTxnResult result = apply(request, OpCode.error, new ErrorTxn(Code.SESSIONEXPIRED.intValue()));
+ request.setException(KeeperException.create(Code.NOAUTH));
+ AuditHelper.addAuditLog(request, result, true);
+ assertWrite(fields(capture.read(1).get(0)), "create", "/denied", "7", "failed", "-102");
+ assertNull(tree.getNode("/denied"));
+ }
+
+ @Test
+ public void testFailedMultiNeverCommitsZeroCodeErrorMembers() throws Exception {
+ Request request = request(OpCode.multi, new MultiOperationRecord(Arrays.asList(
+ Op.create("/rolled", new byte[1], ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT),
+ Op.check("/", -1),
+ Op.setData("/missing", new byte[3], -1),
+ Op.create("/later", new byte[2], ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT))));
+ ProcessTxnResult result = apply(request, OpCode.multi, new MultiTxn(Arrays.asList(
+ txn(OpCode.create, createTxn("/rolled", new byte[1], false)),
+ txn(OpCode.check, new CheckVersionTxn("/", -1)),
+ txn(OpCode.error, new ErrorTxn(Code.NONODE.intValue())),
+ txn(OpCode.error, new ErrorTxn(Code.RUNTIMEINCONSISTENCY.intValue())))));
+ assertEquals(-101, result.err);
+ assertEquals(OpCode.error, result.multiResult.get(0).type);
+ assertEquals(0, result.multiResult.get(0).err);
+ assertNull(tree.getNode("/rolled"));
+ assertNull(tree.getNode("/later"));
+
+ AuditHelper.addAuditLog(request, result);
+ List logs = capture.read(4);
+ assertWrite(fields(logs.get(0)), "multiOperation", null, null, "failed", "-101");
+ assertWrite(fields(logs.get(1)), "create", "/rolled", "1", "rolled_back", "0");
+ assertWrite(fields(logs.get(2)), "setData", "/missing", "3", "failed", "-101");
+ assertWrite(fields(logs.get(3)), "create", "/later", "2", "rolled_back", "-2");
+ assertEquals("0", fields(logs.get(1)).get("multi_index"));
+ assertEquals("2", fields(logs.get(2)).get("multi_index"));
+ assertEquals("3", fields(logs.get(3)).get("multi_index"));
+ for (String log : logs) {
+ assertFalse(log, log.contains("result=success"));
+ assertFalse(log, log.contains("outcome=committed"));
+ }
+ }
+
+ @Test
+ public void testMultiMatchesRepeatedPathsAndTtlByPosition() throws Exception {
+ Request request = request(OpCode.multi, new MultiOperationRecord(Arrays.asList(
+ Op.create("/same", new byte[1], ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT),
+ Op.delete("/same", -1),
+ Op.check("/", -1),
+ Op.create("/same", new byte[2], ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.EPHEMERAL),
+ Op.create("/seq-", new byte[3], ZooDefs.Ids.OPEN_ACL_UNSAFE,
+ CreateMode.PERSISTENT_SEQUENTIAL_WITH_TTL, 60000))));
+ ProcessTxnResult result = apply(request, OpCode.multi, new MultiTxn(Arrays.asList(
+ txn(OpCode.create, createTxn("/same", new byte[1], false)),
+ txn(OpCode.delete, new DeleteTxn("/same")),
+ txn(OpCode.check, new CheckVersionTxn("/", -1)),
+ txn(OpCode.create, createTxn("/same", new byte[2], true)),
+ txn(OpCode.createTTL, new CreateTTLTxn(
+ "/seq-0000000002", new byte[3], ZooDefs.Ids.OPEN_ACL_UNSAFE, -1, 60000)))));
+ assertEquals(0, result.err);
+ assertNotNull(tree.getNode("/seq-0000000002"));
+
+ AuditHelper.addAuditLog(request, result);
+ List logs = capture.read(4);
+ assertWrite(fields(logs.get(0)), "create", "/same", "1", "committed", "0");
+ assertEquals("persistent", fields(logs.get(0)).get("znode_type"));
+ assertEquals("ephemeral", fields(logs.get(2)).get("znode_type"));
+ assertEquals("3", fields(logs.get(2)).get("multi_index"));
+ assertWrite(fields(logs.get(3)), "create", "/seq-0000000002", "3", "committed", "0");
+ assertEquals("persistent_sequential_with_ttl", fields(logs.get(3)).get("znode_type"));
+ assertEquals("4", fields(logs.get(3)).get("multi_index"));
+ }
+
+ @Test
+ public void testIncompleteMultiResultsDoNotMisattributeCodes() throws Exception {
+ Request request = request(OpCode.multi, new MultiOperationRecord(Arrays.asList(
+ Op.create("/incomplete", new byte[2], ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT),
+ Op.check("/", -1))));
+ ProcessTxnResult result = apply(request, OpCode.multi, new MultiTxn(Arrays.asList(
+ txn(OpCode.create, createTxn("/incomplete", new byte[2], false)),
+ txn(OpCode.check, new CheckVersionTxn("/", -1)))));
+ assertEquals(0, result.err);
+ result.multiResult.remove(0);
+ long before = auditErrors();
+ AuditHelper.addAuditLog(request, result);
+ assertWrite(fields(capture.read(1).get(0)), "create", "/incomplete", "2", "unknown", null);
+ assertEquals(before + 1, auditErrors());
+ }
+
+ @Test
+ public void testFirstRuntimeInconsistencyIsFailureNotSkippedMember() throws Exception {
+ Request request = request(OpCode.multi, new MultiOperationRecord(Arrays.asList(
+ Op.create("/first", new byte[1], ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT),
+ Op.create("/skipped", new byte[2], ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT))));
+ ProcessTxnResult result = apply(request, OpCode.multi, new MultiTxn(Arrays.asList(
+ txn(OpCode.error, new ErrorTxn(Code.RUNTIMEINCONSISTENCY.intValue())),
+ txn(OpCode.error, new ErrorTxn(Code.RUNTIMEINCONSISTENCY.intValue())))));
+ assertNull(tree.getNode("/first"));
+ assertNull(tree.getNode("/skipped"));
+ AuditHelper.addAuditLog(request, result);
+ List logs = capture.read(3);
+ assertWrite(fields(logs.get(1)), "create", "/first", "1", "failed", "-2");
+ assertWrite(fields(logs.get(2)), "create", "/skipped", "2", "rolled_back", "-2");
+ }
+
+ @Test
+ public void testDecodingPreservesPositionLimitAndMark() throws Exception {
+ for (int kind = 0; kind < 3; kind++) {
+ String path = "/buffer-" + kind;
+ byte[] bytes = serialize(createRecord(path, new byte[3], CreateMode.PERSISTENT));
+ ByteBuffer storage = kind == 0
+ ? ByteBuffer.allocateDirect(bytes.length + 12) : ByteBuffer.allocate(bytes.length + 12);
+ storage.position(4);
+ storage.put(bytes);
+ storage.position(4);
+ ByteBuffer buffer = storage.slice();
+ buffer.limit(bytes.length);
+ if (kind == 2) {
+ buffer = buffer.asReadOnlyBuffer();
+ }
+ buffer.position(2);
+ buffer.mark();
+ buffer.position(bytes.length);
+ Request request = request(OpCode.create, buffer);
+ ProcessTxnResult result = apply(request, OpCode.create, createTxn(path, new byte[3], false));
+ AuditHelper.addAuditLog(request, result);
+ assertEquals(bytes.length, buffer.position());
+ assertEquals(bytes.length, buffer.limit());
+ buffer.reset();
+ assertEquals(2, buffer.position());
+ assertWrite(fields(capture.read(1).get(0)), "create", path, "3", "committed", "0");
+ }
+ }
+
+ @Test
+ public void testMultiDecodingPreservesPositionLimitAndMark() throws Exception {
+ Request request = request(OpCode.multi, new MultiOperationRecord(Collections.singletonList(
+ Op.create("/buffer-multi", new byte[2], ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT))));
+ request.request.position(3);
+ request.request.mark();
+ int limit = request.request.limit();
+ request.request.position(limit);
+ ProcessTxnResult result = apply(request, OpCode.multi, new MultiTxn(Collections.singletonList(
+ txn(OpCode.create, createTxn("/buffer-multi", new byte[2], false)))));
+ AuditHelper.addAuditLog(request, result);
+ assertEquals(limit, request.request.position());
+ assertEquals(limit, request.request.limit());
+ request.request.reset();
+ assertEquals(3, request.request.position());
+ assertWrite(fields(capture.read(1).get(0)), "create", "/buffer-multi", "2", "committed", "0");
+ }
+
+ @Test
+ public void testUnavailablePayloadIsOmittedAndDecodeErrorsAreCounted() throws Exception {
+ for (ByteBuffer buffer : Arrays.asList(null, ByteBuffer.wrap(new byte[] {0, 0, 0, 10, 1}))) {
+ String path = buffer == null ? "/unavailable" : "/undecodable";
+ Request request = request(OpCode.create, buffer);
+ ProcessTxnResult result = apply(request, OpCode.create, createTxn(path, new byte[5], false));
+ long before = auditErrors();
+ AuditHelper.addAuditLog(request, result);
+ assertWrite(fields(capture.read(1).get(0)), "create", path, null, "committed", "0");
+ assertEquals(before + 1, auditErrors());
+ assertNotNull(tree.getNode(path));
+ }
+ }
+
+ @Test
+ public void testFailedMultiParentSurvivesUndecodableRequest() throws Exception {
+ Request request = request(OpCode.multi, ByteBuffer.wrap(new byte[] {1}));
+ ProcessTxnResult result = apply(request, OpCode.multi, new MultiTxn(Collections.singletonList(
+ txn(OpCode.error, new ErrorTxn(Code.NOAUTH.intValue())))));
+ long before = auditErrors();
+ AuditHelper.addAuditLog(request, result);
+ assertWrite(fields(capture.read(1).get(0)), "multiOperation", null, null, "failed", "-102");
+ assertEquals(before + 1, auditErrors());
+ }
+
+ @Test
+ public void testMissingResultDoesNotInventSuccessOrZxid() throws Exception {
+ Request request = request(OpCode.setData, new SetDataRequest("/unknown", new byte[4], -1));
+ try {
+ AuditHelper.addAuditLog(request, null);
+ } catch (RuntimeException e) {
+ fail("Unavailable transaction results must not escape the audit boundary: " + e);
+ }
+ Map fields = fields(capture.read(1).get(0));
+ assertWrite(fields, "setData", "/unknown", "4", "unknown", null);
+ assertEquals("invoked", fields.get("result"));
+ assertNull(fields.get("zxid"));
+ }
+
+ @Test
+ public void testEnhancedSetAclDoesNotExposeCredentials() throws Exception {
+ List acls = Arrays.asList(
+ new ACL(ZooDefs.Perms.ALL, new Id("world", "anyone")),
+ new ACL(ZooDefs.Perms.READ, new Id("digest", "alice:synthetic-digest-secret")),
+ new ACL(ZooDefs.Perms.WRITE, new Id("custom", "synthetic-token-secret")),
+ new ACL(ZooDefs.Perms.READ, new Id("x509", "-----BEGIN CERTIFICATE-----synthetic-body")));
+ Request request = request(OpCode.setACL, new SetACLRequest("/acl", acls, -1));
+ ProcessTxnResult result = apply(request, OpCode.error, new ErrorTxn(Code.INVALIDACL.intValue()));
+ AuditHelper.addAuditLog(request, result, true);
+ String log = capture.read(1).get(0);
+ assertWrite(fields(log), "setAcl", "/acl", null, "failed", "-114");
+ assertTrue(log, log.contains("world:anyone:cdrwa"));
+ assertTrue(log, log.contains("digest:alice:r"));
+ assertFalse(log, log.contains("synthetic-digest-secret"));
+ assertFalse(log, log.contains("synthetic-token-secret"));
+ assertFalse(log, log.contains("BEGIN CERTIFICATE"));
+ assertFalse(log, log.contains("synthetic-body"));
+ }
+
+ @Test
+ public void testMalformedDigestIdentityIsNotMistakenForAUser() throws Exception {
+ Request request = request(OpCode.setACL, new SetACLRequest("/acl", Collections.singletonList(
+ new ACL(ZooDefs.Perms.READ, new Id("digest", "synthetic-digest-token"))), -1));
+ ProcessTxnResult result = apply(request, OpCode.error, new ErrorTxn(Code.INVALIDACL.intValue()));
+ AuditHelper.addAuditLog(request, result, true);
+ String log = capture.read(1).get(0);
+ assertFalse(log, log.contains("synthetic-digest-token"));
+ assertWrite(fields(log), "setAcl", "/acl", null, "failed", "-114");
+ }
+
+ @Test
+ public void testEnhancedUsersRedactRegisteredDefaultProvider() throws Exception {
+ String property = ProviderRegistry.AUTHPROVIDER_PROPERTY_PREFIX + "c1-audit-user";
+ String previous = System.getProperty(property);
+ System.setProperty(property, CredentialAuthenticationProvider.class.getName());
+ ProviderRegistry.initialize();
+ try {
+ Request request = new Request(cnxn, SESSION, 41, OpCode.create,
+ ByteBuffer.wrap(serialize(createRecord("/custom-user", new byte[1], CreateMode.PERSISTENT))),
+ Arrays.asList(new Id("ip", "127.0.0.1"),
+ new Id("audit-test-custom", "alice:synthetic-password")));
+ ProcessTxnResult result = apply(request, OpCode.create, createTxn("/custom-user", new byte[1], false));
+ assertEquals(0, result.err);
+ AuditHelper.addAuditLog(request, result);
+ String log = capture.read(1).get(0);
+ assertWrite(fields(log), "create", "/custom-user", "1", "committed", "0");
+ assertEquals("127.0.0.1,[redacted]", fields(log).get("user"));
+ assertFalse(log.contains("synthetic-password"));
+ } finally {
+ ProviderRegistry.removeProvider("audit-test-custom");
+ restoreProperty(property, previous);
+ }
+ }
+
+ @Test
+ public void testEnhancedUsersRedactUnknownAndMalformedIdentities() throws Exception {
+ Request request = new Request(cnxn, SESSION, 41, OpCode.create,
+ ByteBuffer.wrap(serialize(createRecord("/unknown-users", new byte[1], CreateMode.PERSISTENT))),
+ Arrays.asList(new Id("unregistered", "synthetic-token"),
+ new Id("digest", "synthetic-digest-token"),
+ new Id("ip", "alice:synthetic-password")));
+ ProcessTxnResult result = apply(request, OpCode.create, createTxn("/unknown-users", new byte[1], false));
+ AuditHelper.addAuditLog(request, result);
+ String log = capture.read(1).get(0);
+ assertWrite(fields(log), "create", "/unknown-users", "1", "committed", "0");
+ assertEquals("[redacted],[redacted],[redacted]", fields(log).get("user"));
+ assertFalse(log.contains("synthetic"));
+ }
+
+ @Test
+ public void testAclModeSnapshotSurvivesOffToOnInterleaving() throws Exception {
+ Request create = request(OpCode.create, createRecord("/mode-acl", new byte[0], CreateMode.PERSISTENT));
+ assertEquals(0, apply(create, OpCode.create, createTxn("/mode-acl", new byte[0], false)).err);
+ String digest = DigestAuthenticationProvider.generateDigest("alice:synthetic-password");
+ List acls = Collections.singletonList(new ACL(ZooDefs.Perms.ALL, new Id("digest", digest)));
+ SetACLRequest record = new SetACLRequest("/mode-acl", acls, -1);
+ AtomicInteger transitions = new AtomicInteger();
+ Request switching = new Request(cnxn, SESSION, 41, OpCode.setACL, ByteBuffer.wrap(serialize(record)),
+ Collections.singletonList(new Id("ip", "127.0.0.1"))) {
+ @Override
+ public String getUsers() {
+ System.setProperty(ENHANCED_ENABLE, "true");
+ transitions.incrementAndGet();
+ return super.getUsers();
+ }
+ };
+ ProcessTxnResult changed = apply(switching, OpCode.setACL, new SetACLTxn("/mode-acl", acls, 1));
+ assertEquals(0, changed.err);
+ System.setProperty(ENHANCED_ENABLE, "false");
+ AuditHelper.addAuditLog(switching, changed);
+ Map legacy = fields(capture.read(1).get(0));
+ assertEquals(1, transitions.get());
+ assertEquals("true", System.getProperty(ENHANCED_ENABLE));
+ assertNull("An in-flight legacy record must not be relabeled v2", legacy.get("schema_version"));
+ assertEquals("digest:" + digest + ":cdrwa", legacy.get("acl"));
+ assertEquals("success", legacy.get("result"));
+
+ Request enhanced = request(OpCode.setACL, record);
+ ProcessTxnResult updated = apply(enhanced, OpCode.setACL, new SetACLTxn("/mode-acl", acls, 2));
+ assertEquals(0, updated.err);
+ AuditHelper.addAuditLog(enhanced, updated);
+ String log = capture.read(1).get(0);
+ assertWrite(fields(log), "setAcl", "/mode-acl", null, "committed", "0");
+ assertEquals("digest:alice:cdrwa", fields(log).get("acl"));
+ assertFalse(log.contains(digest));
+ assertFalse(log.contains("synthetic-password"));
+ }
+
+ @Test
+ public void testSuccessfulMultiKeepsOneModeAcrossMembers() throws Exception {
+ Request request = request(OpCode.multi, new MultiOperationRecord(Arrays.asList(
+ Op.create("/mode-multi", new byte[1], ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT),
+ Op.setData("/mode-multi", new byte[2], -1))));
+ ProcessTxnResult result = apply(request, OpCode.multi, new MultiTxn(Arrays.asList(
+ txn(OpCode.create, createTxn("/mode-multi", new byte[1], false)),
+ txn(OpCode.setData, new SetDataTxn("/mode-multi", new byte[2], 1)))));
+ assertEquals(0, result.err);
+ System.setProperty(ENHANCED_ENABLE, "false");
+ AtomicInteger events = new AtomicInteger();
+ AuditLogger delegate = new Slf4jAuditLogger();
+ Object previousLogger = replaceProviderField("auditLogger", (AuditLogger) event -> {
+ delegate.logAuditEvent(event);
+ if (events.incrementAndGet() == 1) {
+ System.setProperty(ENHANCED_ENABLE, "true");
+ }
+ });
+ try {
+ AuditHelper.addAuditLog(request, result);
+ List logs = capture.read(2);
+ assertEquals(2, events.get());
+ for (String log : logs) {
+ assertNull("A multi must keep its captured legacy mode", fields(log).get("schema_version"));
+ assertEquals("success", fields(log).get("result"));
+ }
+
+ Request following = request(OpCode.setData, new SetDataRequest("/mode-multi", new byte[3], -1));
+ ProcessTxnResult updated = apply(following, OpCode.setData, new SetDataTxn("/mode-multi", new byte[3], 2));
+ assertEquals(0, updated.err);
+ AuditHelper.addAuditLog(following, updated);
+ assertWrite(fields(capture.read(1).get(0)), "setData", "/mode-multi", "3", "committed", "0");
+ } finally {
+ replaceProviderField("auditLogger", previousLogger);
+ }
+ }
+
+ @Test
+ public void testFailedMultiKeepsParentModeForRolledBackMembers() throws Exception {
+ Request request = request(OpCode.multi, new MultiOperationRecord(Arrays.asList(
+ Op.create("/mode-rolled", new byte[1], ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT),
+ Op.check("/missing", -1),
+ Op.setData("/mode-rolled", new byte[2], -1))));
+ ProcessTxnResult result = apply(request, OpCode.multi, new MultiTxn(Arrays.asList(
+ txn(OpCode.create, createTxn("/mode-rolled", new byte[1], false)),
+ txn(OpCode.error, new ErrorTxn(Code.NONODE.intValue())),
+ txn(OpCode.error, new ErrorTxn(Code.RUNTIMEINCONSISTENCY.intValue())))));
+ assertEquals(-101, result.err);
+ assertNull(tree.getNode("/mode-rolled"));
+ AtomicInteger events = new AtomicInteger();
+ AuditLogger delegate = new Slf4jAuditLogger();
+ Object previousLogger = replaceProviderField("auditLogger", (AuditLogger) event -> {
+ delegate.logAuditEvent(event);
+ if (events.incrementAndGet() == 1) {
+ System.setProperty(ENHANCED_ENABLE, "false");
+ }
+ });
+ try {
+ AuditHelper.addAuditLog(request, result);
+ List logs = capture.read(3);
+ assertEquals(3, events.get());
+ assertWrite(fields(logs.get(0)), "multiOperation", null, null, "failed", "-101");
+ assertWrite(fields(logs.get(1)), "create", "/mode-rolled", "1", "rolled_back", "0");
+ assertWrite(fields(logs.get(2)), "setData", "/mode-rolled", "2", "rolled_back", "-2");
+ assertEquals("2", fields(logs.get(2)).get("multi_index"));
+ } finally {
+ replaceProviderField("auditLogger", previousLogger);
+ }
+ }
+
+ @Test
+ public void testAuditDisabledSkipsRequestsAndProvider() throws Exception {
+ Object previous = replaceProviderField("auditEnabled", false);
+ long before = auditErrors();
+ try {
+ AuditHelper.addAuditLog(null, null);
+ ZKAuditProvider.log("user", "create", "/disabled", null, null, null, null, Result.SUCCESS);
+ capture.read(0);
+ assertEquals(before, auditErrors());
+ } finally {
+ replaceProviderField("auditEnabled", previous);
+ }
+ }
+
+ @Test
+ public void testDefaultAndExplicitLegacyOutputAreUnchanged() throws Exception {
+ for (String setting : Arrays.asList(null, "false")) {
+ restoreProperty(ENHANCED_ENABLE, setting);
+ ZKAuditProvider.log("user", "create", "/legacy", null, "persistent", "0x123",
+ "127.0.0.1", Result.SUCCESS);
+ assertEquals("session=0x123\tuser=user\tip=127.0.0.1\toperation=create"
+ + "\tznode=/legacy\tznode_type=persistent\tresult=success", capture.read(1).get(0));
+ }
+ }
+
+ @Test
+ public void testLegacyFailedMultiStillHasOnlyParent() throws Exception {
+ System.clearProperty(ENHANCED_ENABLE);
+ Request request = request(OpCode.multi, new MultiOperationRecord(Collections.singletonList(
+ Op.setData("/absent", new byte[2], -1))));
+ ProcessTxnResult result = apply(request, OpCode.multi, new MultiTxn(Collections.singletonList(
+ txn(OpCode.error, new ErrorTxn(Code.NONODE.intValue())))));
+ AuditHelper.addAuditLog(request, result);
+ Map fields = fields(capture.read(1).get(0));
+ assertEquals("multiOperation", fields.get("operation"));
+ assertEquals("failure", fields.get("result"));
+ assertNull(fields.get("schema_version"));
+ }
+
+ @Test
+ public void testReadsRemainUnaudited() {
+ for (int type : new int[] {OpCode.getData, OpCode.getChildren, OpCode.exists, OpCode.multiRead}) {
+ AuditHelper.addAuditLog(request(type, ByteBuffer.wrap(new byte[] {1})), new ProcessTxnResult());
+ }
+ capture.read(0);
+ }
+
+ @Test
+ public void testSystemDeletionHasStableIdentityAndSystemActor() throws Exception {
+ Request create = request(OpCode.create, createRecord("/ephemeral", new byte[0], CreateMode.EPHEMERAL));
+ assertEquals(0, apply(create, OpCode.create, createTxn("/ephemeral", new byte[0], true)).err);
+ tree.processTxn(new TxnHeader(SESSION, -11, 80, 1000, OpCode.closeSession),
+ new CloseSessionTxn(Collections.singletonList("/ephemeral")));
+ Map fields = fields(capture.read(1).get(0));
+ assertWrite(fields, "ephemeralZNodeDeletionOnSessionCloseOrExpire",
+ "/ephemeral", null, "committed", "0");
+ assertEquals(ZKAuditProvider.getZKUser(), fields.get("user"));
+ assertEquals("0x123", fields.get("session"));
+ assertEquals("80", fields.get("zxid"));
+ assertNull(fields.get("cxid"));
+ assertNull(fields.get("ip"));
+ assertNull(tree.getNode("/ephemeral"));
+ }
+
+ @Test
+ public void testLoggerFailureDoesNotInterruptSystemDeletion() throws Exception {
+ for (String path : Arrays.asList("/ephemeral-a", "/ephemeral-b")) {
+ Request request = request(OpCode.create, createRecord(path, new byte[0], CreateMode.EPHEMERAL));
+ assertEquals(0, apply(request, OpCode.create, createTxn(path, new byte[0], true)).err);
+ }
+ AuditLogger failingLogger = event -> {
+ throw new IllegalStateException("synthetic audit sink failure");
+ };
+ Object previous = replaceProviderField("auditLogger", failingLogger);
+ long before = auditErrors();
+ try {
+ try {
+ tree.processTxn(new TxnHeader(SESSION, -11, 81, 1000, OpCode.closeSession),
+ new CloseSessionTxn(Arrays.asList("/ephemeral-a", "/ephemeral-b")));
+ } catch (RuntimeException e) {
+ fail("Audit sink failure must not interrupt committed deletion: " + e);
+ }
+ assertNull(tree.getNode("/ephemeral-a"));
+ assertNull(tree.getNode("/ephemeral-b"));
+ assertEquals(before + 2, auditErrors());
+ } finally {
+ replaceProviderField("auditLogger", previous);
+ }
+ }
+
+ @Test
+ public void testFailingErrorCounterDoesNotEscapeMetadataFailure() throws Exception {
+ Request request = request(OpCode.create, ByteBuffer.wrap(new byte[] {1}));
+ ProcessTxnResult result = apply(request, OpCode.create, createTxn("/reporter-metadata", new byte[2], false));
+ FailingCounter counter = new FailingCounter();
+ Counter previousCounter = replaceAuditErrorCounter(counter);
+ try {
+ RuntimeException escaped = null;
+ try {
+ AuditHelper.addAuditLog(request, result);
+ } catch (RuntimeException e) {
+ escaped = e;
+ }
+ assertNotNull(tree.getNode("/reporter-metadata"));
+ assertNull("The error reporter must not escape the audit boundary", escaped);
+ assertWrite(fields(capture.read(1).get(0)), "create", "/reporter-metadata", null, "committed", "0");
+ assertEquals("Do not retry a failing error reporter", 1, counter.get());
+ } finally {
+ replaceAuditErrorCounter(previousCounter);
+ }
+ }
+
+ @Test
+ public void testFailingErrorCounterDoesNotInterruptSystemDeletions() throws Exception {
+ for (String path : Arrays.asList("/reporter-a", "/reporter-b")) {
+ Request request = request(OpCode.create, createRecord(path, new byte[0], CreateMode.EPHEMERAL));
+ assertEquals(0, apply(request, OpCode.create, createTxn(path, new byte[0], true)).err);
+ }
+ FailingCounter counter = new FailingCounter();
+ Counter previousCounter = replaceAuditErrorCounter(counter);
+ Object previousLogger = replaceProviderField("auditLogger", (AuditLogger) event -> {
+ throw new IllegalStateException("synthetic audit sink failure");
+ });
+ try {
+ RuntimeException escaped = null;
+ try {
+ tree.processTxn(new TxnHeader(SESSION, -11, 81, 1000, OpCode.closeSession),
+ new CloseSessionTxn(Arrays.asList("/reporter-a", "/reporter-b")));
+ } catch (RuntimeException e) {
+ escaped = e;
+ }
+ assertNull("Reporting a sink failure must not escape system deletion", escaped);
+ assertNull(tree.getNode("/reporter-a"));
+ assertNull(tree.getNode("/reporter-b"));
+ assertEquals("One best-effort report per deletion, with no retries", 2, counter.get());
+ } finally {
+ replaceProviderField("auditLogger", previousLogger);
+ replaceAuditErrorCounter(previousCounter);
+ }
+ }
+
+ private Request request(int type, Record record) throws IOException {
+ return request(type, ByteBuffer.wrap(serialize(record)));
+ }
+
+ private Request request(int type, ByteBuffer buffer) {
+ return new Request(cnxn, SESSION, 41, type, buffer,
+ Collections.singletonList(new Id("ip", "127.0.0.1")));
+ }
+
+ private ProcessTxnResult apply(Request request, int type, Record txn) {
+ TxnHeader header = new TxnHeader(SESSION, 41, 66, 1000, type);
+ request.setHdr(header);
+ request.setTxn(txn);
+ return tree.processTxn(header, txn);
+ }
+
+ private static CreateRequest createRecord(String path, byte[] data, CreateMode mode) {
+ return new CreateRequest(path, data, ZooDefs.Ids.OPEN_ACL_UNSAFE, mode.toFlag());
+ }
+
+ private static CreateTxn createTxn(String path, byte[] data, boolean ephemeral) {
+ return new CreateTxn(path, data, ZooDefs.Ids.OPEN_ACL_UNSAFE, ephemeral, -1);
+ }
+
+ private static Txn txn(int type, Record record) throws IOException {
+ return new Txn(type, serialize(record));
+ }
+
+ private static byte[] serialize(Record record) throws IOException {
+ ByteArrayOutputStream out = new ByteArrayOutputStream();
+ record.serialize(BinaryOutputArchive.getArchive(out), "request");
+ return out.toByteArray();
+ }
+
+ static long auditErrors() {
+ return ((Number) MetricsUtils.currentServerMetrics().getOrDefault("audit_errors", 0L)).longValue();
+ }
+
+ static Object replaceProviderField(String name, Object value) throws ReflectiveOperationException {
+ Field field = ZKAuditProvider.class.getDeclaredField(name);
+ field.setAccessible(true);
+ Object previous = field.get(null);
+ field.set(null, value);
+ return previous;
+ }
+
+ static Counter replaceAuditErrorCounter(Counter counter) throws ReflectiveOperationException {
+ ServerMetrics metrics = ServerMetrics.getMetrics();
+ Field field = ServerMetrics.class.getField("AUDIT_ERRORS");
+ field.setAccessible(true);
+ Counter previous = (Counter) field.get(metrics);
+ field.set(metrics, counter);
+ return previous;
+ }
+
+ static void restoreProperty(String name, String value) {
+ if (value == null) {
+ System.clearProperty(name);
+ } else {
+ System.setProperty(name, value);
+ }
+ }
+
+ static Map fields(String log) {
+ Map fields = new LinkedHashMap<>();
+ for (String pair : log.split("\t")) {
+ int separator = pair.indexOf('=');
+ assertTrue(log, separator > 0);
+ String previous = fields.put(pair.substring(0, separator), pair.substring(separator + 1));
+ assertNull("Duplicate audit field in " + log, previous);
+ }
+ return fields;
+ }
+
+ static void assertWrite(Map fields, String operation, String path,
+ String length, String outcome, String error) {
+ assertEquals("2", fields.get("schema_version"));
+ assertEquals(operation, fields.get("operation"));
+ assertEquals(path, fields.get("znode"));
+ assertEquals(length, fields.get("data_length"));
+ assertEquals(outcome, fields.get("outcome"));
+ assertEquals(error, fields.get("error_code"));
+ assertEquals("committed".equals(outcome) ? "success" : "unknown".equals(outcome) ? "invoked" : "failure",
+ fields.get("result"));
+ }
+
+ public static class CredentialAuthenticationProvider implements AuthenticationProvider {
+ @Override
+ public String getScheme() {
+ return "audit-test-custom";
+ }
+
+ @Override
+ public Code handleAuthentication(ServerCnxn connection, byte[] authData) {
+ connection.addAuthInfo(new Id(getScheme(), new String(authData, StandardCharsets.UTF_8)));
+ return Code.OK;
+ }
+
+ @Override
+ public boolean matches(String id, String aclExpr) {
+ return id.equals(aclExpr);
+ }
+
+ @Override
+ public boolean isAuthenticated() {
+ return true;
+ }
+
+ @Override
+ public boolean isValid(String id) {
+ return true;
+ }
+ }
+
+ static final class FailingCounter implements Counter {
+ private final AtomicInteger attempts = new AtomicInteger();
+
+ @Override
+ public void add(long delta) {
+ attempts.incrementAndGet();
+ throw new IllegalStateException("synthetic audit counter failure");
+ }
+
+ @Override
+ public long get() {
+ return attempts.get();
+ }
+ }
+
+ static final class AuditCapture extends AppenderBase implements AutoCloseable {
+ private final Logger logger = (Logger) LoggerFactory.getLogger(Slf4jAuditLogger.class);
+ private final Level previousLevel = logger.getLevel();
+ private final List messages = new ArrayList<>();
+ private boolean overflow;
+
+ AuditCapture() {
+ setContext(logger.getLoggerContext());
+ logger.setLevel(Level.INFO);
+ logger.addAppender(this);
+ start();
+ }
+
+ @Override
+ protected synchronized void append(ILoggingEvent event) {
+ if (messages.size() < 128) {
+ messages.add(event.getFormattedMessage());
+ } else {
+ overflow = true;
+ }
+ notifyAll();
+ }
+
+ synchronized List read(int expected) {
+ assertFalse("Audit capture exceeded its bounded capacity", overflow);
+ List result = new ArrayList<>(messages);
+ messages.clear();
+ assertEquals(result.toString(), expected, result.size());
+ return result;
+ }
+
+ synchronized void clear() {
+ assertFalse("Audit capture exceeded its bounded capacity", overflow);
+ messages.clear();
+ }
+
+ synchronized List await(int expected, long timeoutMillis) throws InterruptedException {
+ long deadline = System.nanoTime() + TimeUnit.MILLISECONDS.toNanos(timeoutMillis);
+ while (messages.size() < expected) {
+ long remaining = deadline - System.nanoTime();
+ if (remaining <= 0) {
+ break;
+ }
+ TimeUnit.NANOSECONDS.timedWait(this, remaining);
+ }
+ return read(expected);
+ }
+
+ @Override
+ public void close() {
+ logger.detachAppender(this);
+ logger.setLevel(previousLevel);
+ stop();
+ }
+ }
+}
diff --git a/zookeeper-server/src/test/java/org/apache/zookeeper/audit/Slf4JAuditLoggerTest.java b/zookeeper-server/src/test/java/org/apache/zookeeper/audit/Slf4JAuditLoggerTest.java
index f8d3c7b47af..0f426e073f7 100644
--- a/zookeeper-server/src/test/java/org/apache/zookeeper/audit/Slf4JAuditLoggerTest.java
+++ b/zookeeper-server/src/test/java/org/apache/zookeeper/audit/Slf4JAuditLoggerTest.java
@@ -19,14 +19,12 @@
import static org.apache.zookeeper.test.ClientBase.CONNECTION_TIMEOUT;
import static org.junit.Assert.assertEquals;
-import java.io.ByteArrayOutputStream;
import java.io.IOException;
-import java.io.LineNumberReader;
-import java.io.StringReader;
import java.net.InetAddress;
import java.net.InetSocketAddress;
import java.util.ArrayList;
import java.util.List;
+import java.util.Map;
import org.apache.zookeeper.CreateMode;
import org.apache.zookeeper.KeeperException;
import org.apache.zookeeper.KeeperException.Code;
@@ -36,6 +34,7 @@
import org.apache.zookeeper.ZooDefs;
import org.apache.zookeeper.ZooKeeper;
import org.apache.zookeeper.audit.AuditEvent.Result;
+import org.apache.zookeeper.audit.AuditHelperTest.AuditCapture;
import org.apache.zookeeper.data.ACL;
import org.apache.zookeeper.data.Stat;
import org.apache.zookeeper.server.Request;
@@ -43,7 +42,6 @@
import org.apache.zookeeper.server.quorum.QuorumPeerTestBase;
import org.apache.zookeeper.test.ClientBase;
import org.apache.zookeeper.test.ClientBase.CountdownWatcher;
-import org.apache.zookeeper.test.LoggerTestTool;
import org.junit.AfterClass;
import org.junit.Assert;
import org.junit.Before;
@@ -57,14 +55,14 @@ public class Slf4JAuditLoggerTest extends QuorumPeerTestBase {
private static int SERVER_COUNT = 3;
private static MainThread[] mt;
private static ZooKeeper zk;
- private static ByteArrayOutputStream os;
+ private static AuditCapture os;
@BeforeClass
public static void setUpBeforeClass() throws Exception {
System.setProperty(ZKAuditProvider.AUDIT_ENABLE, "true");
+ System.setProperty("zookeeper.extendedTypesEnabled", "true");
// setup the logger to capture all logs
- LoggerTestTool loggerTestTool = new LoggerTestTool(Slf4jAuditLogger.class);
- os = loggerTestTool.getOutputStream();
+ os = new AuditCapture();
mt = startQuorum();
zk = ClientBase.createZKClient("127.0.0.1:" + mt[0].getQuorumPeer().getClientPort());
//Verify start audit log here itself
@@ -75,7 +73,7 @@ public static void setUpBeforeClass() throws Exception {
@Before
public void setUp() {
- os.reset();
+ os.clear();
}
@Test
@@ -102,13 +100,47 @@ public void testCreateAuditLogs()
null, createMode), readAuditLog(os));
}
+ @Test
+ public void testCreateWithTtlAuditLogs() throws Exception {
+ String path = zk.create("/createTtlPath", new byte[0], ZooDefs.Ids.OPEN_ACL_UNSAFE,
+ CreateMode.PERSISTENT_WITH_TTL, null, 60000);
+ verifyLog(getAuditLog(AuditConstants.OP_CREATE, path, Result.SUCCESS,
+ null, "persistent_with_ttl"), readAuditLog(os));
+ }
+
+ @Test
+ public void testCreateSequentialWithTtlAuditLogs() throws Exception {
+ String path = zk.create("/createTtlSeqPath", new byte[0], ZooDefs.Ids.OPEN_ACL_UNSAFE,
+ CreateMode.PERSISTENT_SEQUENTIAL_WITH_TTL, null, 60000);
+ verifyLog(getAuditLog(AuditConstants.OP_CREATE, path, Result.SUCCESS,
+ null, "persistent_sequential_with_ttl"), readAuditLog(os));
+ }
+
+ @Test
+ public void testEnhancedEscapingThroughSlf4j() {
+ String previous = System.getProperty(AuditHelperTest.ENHANCED_ENABLE);
+ try (AuditHelperTest.AuditCapture capture = new AuditHelperTest.AuditCapture()) {
+ System.setProperty(AuditHelperTest.ENHANCED_ENABLE, "true");
+ ZKAuditProvider.log("team\tname\r\n\\t", "create", "/name=value\\child",
+ null, "persistent", "0x123", "127.0.0.1", Result.SUCCESS);
+ String log = capture.read(1).get(0);
+ assertEquals("2", AuditHelperTest.fields(log).get("schema_version"));
+ assertEquals("team\\tname\\r\\n\\\\t", AuditHelperTest.fields(log).get("user"));
+ assertEquals("/name=value\\\\child", AuditHelperTest.fields(log).get("znode"));
+ Assert.assertFalse(log.contains("\n"));
+ Assert.assertFalse(log.contains("\r"));
+ } finally {
+ AuditHelperTest.restoreProperty(AuditHelperTest.ENHANCED_ENABLE, previous);
+ }
+ }
+
@Test
public void testDeleteAuditLogs()
throws InterruptedException, IOException, KeeperException {
String path = "/deletePath";
zk.create(path, "".getBytes(), ZooDefs.Ids.OPEN_ACL_UNSAFE,
CreateMode.PERSISTENT);
- os.reset();
+ os.clear();
try {
zk.delete(path, -100);
} catch (KeeperException exception) {
@@ -129,7 +161,7 @@ public void testSetDataAuditLogs()
String path = "/setDataPath";
zk.create(path, "".getBytes(), ZooDefs.Ids.OPEN_ACL_UNSAFE,
CreateMode.PERSISTENT);
- os.reset();
+ os.clear();
try {
zk.setData(path, "newData".getBytes(), -100);
} catch (KeeperException exception) {
@@ -151,7 +183,7 @@ public void testSetACLAuditLogs()
String path = "/aclPath";
zk.create(path, "".getBytes(), ZooDefs.Ids.OPEN_ACL_UNSAFE,
CreateMode.PERSISTENT);
- os.reset();
+ os.clear();
try {
zk.setACL(path, openAclUnsafe, -100);
} catch (KeeperException exception) {
@@ -243,6 +275,32 @@ public void testEphemralZNodeAuditLogs()
ZKAuditProvider.getZKUser(), null), readAuditLog(os, SERVER_COUNT));
}
+ @Test
+ public void testEnhancedSystemDeletionIdentityAcrossReplicas() throws Exception {
+ String previous = System.getProperty(AuditHelperTest.ENHANCED_ENABLE);
+ System.setProperty(AuditHelperTest.ENHANCED_ENABLE, "true");
+ try (ZooKeeper client = ClientBase.createZKClient("127.0.0.1:" + mt[0].getQuorumPeer().getClientPort())) {
+ client.create("/enhanced-ephemeral", new byte[1], ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.EPHEMERAL);
+ String session = "0x" + Long.toHexString(client.getSessionId());
+ os.read(1);
+ client.close();
+ List logs = os.await(SERVER_COUNT, CONNECTION_TIMEOUT);
+ String zxid = Long.toString(zk.exists("/", false).getPzxid());
+ for (String log : logs) {
+ Map fields = AuditHelperTest.fields(log);
+ AuditHelperTest.assertWrite(fields, AuditConstants.OP_DEL_EZNODE_EXP,
+ "/enhanced-ephemeral", null, "committed", "0");
+ assertEquals(zxid, fields.get("zxid"));
+ assertEquals(session, fields.get("session"));
+ assertEquals(ZKAuditProvider.getZKUser(), fields.get("user"));
+ Assert.assertNull(fields.get("cxid"));
+ Assert.assertNull(fields.get("ip"));
+ }
+ } finally {
+ AuditHelperTest.restoreProperty(AuditHelperTest.ENHANCED_ENABLE, previous);
+ }
+ }
+
private static String getStartLog() {
// user=userName operation=ZooKeeperServer start result=success
@@ -309,10 +367,7 @@ private ServerCnxn getServerCnxn() {
}
private static void verifyLog(String expectedLog, String log) {
- String searchString = " - ";
- int logStartIndex = log.indexOf(searchString);
- String auditLog = log.substring(logStartIndex + searchString.length());
- Assert.assertTrue(auditLog.endsWith(expectedLog));
+ Assert.assertTrue(log, log.endsWith(expectedLog));
}
private static void verifyLogs(String expectedLog, List logs) {
@@ -321,37 +376,14 @@ private static void verifyLogs(String expectedLog, List logs) {
}
}
- private String readAuditLog(ByteArrayOutputStream os) throws IOException {
+ private String readAuditLog(AuditCapture os) throws IOException {
return readAuditLog(os, 1).get(0);
}
- private static List readAuditLog(ByteArrayOutputStream os,
+ private static List readAuditLog(AuditCapture os,
int numberOfLogEntry)
throws IOException {
- return readAuditLog(os, numberOfLogEntry, false);
- }
-
- private static List readAuditLog(ByteArrayOutputStream os,
- int numberOfLogEntry,
- boolean skipEphemralDeletion) throws IOException {
- List logs = new ArrayList<>();
- LineNumberReader r = new LineNumberReader(
- new StringReader(os.toString()));
- String line;
- while ((line = r.readLine()) != null) {
- if (skipEphemralDeletion
- && line.contains(AuditConstants.OP_DEL_EZNODE_EXP)) {
- continue;
- }
- logs.add(line);
- }
- os.reset();
- assertEquals(
- "Expected number of log entries are not generated. Logs are "
- + logs,
- numberOfLogEntry, logs.size());
- return logs;
-
+ return os.read(numberOfLogEntry);
}
private static MainThread[] startQuorum() throws IOException {
@@ -411,6 +443,7 @@ private void waitForDeletion(ZooKeeper zooKeeper, String path)
@AfterClass
public static void tearDownAfterClass() {
System.clearProperty(ZKAuditProvider.AUDIT_ENABLE);
+ System.clearProperty("zookeeper.extendedTypesEnabled");
for (int i = 0; i < SERVER_COUNT; i++) {
try {
if (mt[i] != null) {
@@ -419,11 +452,7 @@ public static void tearDownAfterClass() {
} catch (InterruptedException e) {
e.printStackTrace();
}
- try {
- os.close();
- } catch (IOException e) {
- e.printStackTrace();
- }
}
+ os.close();
}
}
diff --git a/zookeeper-server/src/test/java/org/apache/zookeeper/audit/StandaloneServerAuditTest.java b/zookeeper-server/src/test/java/org/apache/zookeeper/audit/StandaloneServerAuditTest.java
index cf8ca8eea36..851cc1e8f0c 100644
--- a/zookeeper-server/src/test/java/org/apache/zookeeper/audit/StandaloneServerAuditTest.java
+++ b/zookeeper-server/src/test/java/org/apache/zookeeper/audit/StandaloneServerAuditTest.java
@@ -19,34 +19,67 @@
package org.apache.zookeeper.audit;
+import static org.apache.zookeeper.audit.AuditHelperTest.assertWrite;
+import static org.apache.zookeeper.audit.AuditHelperTest.fields;
+import static org.junit.Assert.assertArrayEquals;
import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertNull;
import static org.junit.Assert.assertTrue;
-import java.io.ByteArrayOutputStream;
+import static org.junit.Assert.fail;
import java.io.IOException;
-import java.io.LineNumberReader;
-import java.io.StringReader;
-import java.util.ArrayList;
+import java.nio.charset.StandardCharsets;
+import java.util.Arrays;
+import java.util.Collections;
import java.util.List;
+import java.util.Map;
import org.apache.zookeeper.CreateMode;
import org.apache.zookeeper.KeeperException;
+import org.apache.zookeeper.KeeperException.Code;
+import org.apache.zookeeper.Op;
+import org.apache.zookeeper.OpResult;
import org.apache.zookeeper.ZooDefs;
import org.apache.zookeeper.ZooKeeper;
+import org.apache.zookeeper.audit.AuditHelperTest.AuditCapture;
+import org.apache.zookeeper.audit.AuditHelperTest.CredentialAuthenticationProvider;
+import org.apache.zookeeper.audit.AuditHelperTest.FailingCounter;
+import org.apache.zookeeper.data.ACL;
+import org.apache.zookeeper.data.Id;
+import org.apache.zookeeper.data.Stat;
+import org.apache.zookeeper.metrics.Counter;
+import org.apache.zookeeper.server.auth.DigestAuthenticationProvider;
+import org.apache.zookeeper.server.auth.ProviderRegistry;
import org.apache.zookeeper.test.ClientBase;
-import org.apache.zookeeper.test.LoggerTestTool;
+import org.junit.After;
import org.junit.AfterClass;
+import org.junit.Before;
import org.junit.BeforeClass;
import org.junit.Test;
public class StandaloneServerAuditTest extends ClientBase {
- private static ByteArrayOutputStream os;
+ private AuditCapture capture;
+ private String previousEnhanced;
+ private String previousExtendedTypes;
@BeforeClass
public static void setup() {
System.setProperty(ZKAuditProvider.AUDIT_ENABLE, "true");
- LoggerTestTool loggerTestTool = new LoggerTestTool(Slf4jAuditLogger.class);
- os = loggerTestTool.getOutputStream();
+ }
+
+ @Before
+ public void captureAuditLogs() {
+ previousEnhanced = System.getProperty(AuditHelperTest.ENHANCED_ENABLE);
+ previousExtendedTypes = System.getProperty("zookeeper.extendedTypesEnabled");
+ capture = new AuditCapture();
+ }
+
+ @After
+ public void restoreAuditSettings() {
+ capture.close();
+ AuditHelperTest.restoreProperty(AuditHelperTest.ENHANCED_ENABLE, previousEnhanced);
+ AuditHelperTest.restoreProperty("zookeeper.extendedTypesEnabled", previousExtendedTypes);
}
@AfterClass
@@ -60,21 +93,256 @@ public void testCreateAuditLog() throws KeeperException, InterruptedException, I
String path = "/createPath";
zk.create(path, "".getBytes(), ZooDefs.Ids.OPEN_ACL_UNSAFE,
CreateMode.PERSISTENT);
- List logs = readAuditLog(os);
+ List logs = capture.read(1);
assertEquals(1, logs.size());
assertTrue(logs.get(0).endsWith("operation=create\tznode=/createPath\tznode_type=persistent\tresult=success"));
}
- private static List readAuditLog(ByteArrayOutputStream os) throws IOException {
- List logs = new ArrayList<>();
- LineNumberReader r = new LineNumberReader(
- new StringReader(os.toString()));
- String line;
- while ((line = r.readLine()) != null) {
- logs.add(line);
+ @Test
+ public void testEnhancedWriteResultsMatchClientResponses() throws Exception {
+ System.setProperty(AuditHelperTest.ENHANCED_ENABLE, "true");
+ ZooKeeper zk = createClient();
+ byte[] data = "audit-secret".getBytes(StandardCharsets.UTF_8);
+ zk.create("/enhanced", data, ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT);
+ String log = capture.read(1).get(0);
+ Map created = fields(log);
+ assertWrite(created, "create", "/enhanced", "12", "committed", "0");
+ assertFalse(log.contains("audit-secret"));
+ assertEquals(Long.toString(zk.exists("/enhanced", false).getCzxid()), created.get("zxid"));
+
+ try {
+ zk.create("/enhanced", new byte[7], ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT);
+ fail("Duplicate create must fail");
+ } catch (KeeperException e) {
+ assertEquals(Code.NODEEXISTS, e.code());
}
- os.reset();
- return logs;
+ assertWrite(fields(capture.read(1).get(0)), "create", "/enhanced", "7", "failed", "-110");
+
+ byte[] unicode = "\u00e9\ud83d\ude00".getBytes(StandardCharsets.UTF_8);
+ Stat changed = zk.setData("/enhanced", unicode, -1);
+ Map changedLog = fields(capture.read(1).get(0));
+ assertWrite(changedLog, "setData", "/enhanced", "6", "committed", "0");
+ assertEquals(Long.toString(changed.getMzxid()), changedLog.get("zxid"));
+ try {
+ zk.setData("/enhanced", new byte[1], -100);
+ fail("Invalid version must fail");
+ } catch (KeeperException e) {
+ assertEquals(Code.BADVERSION, e.code());
+ }
+ assertWrite(fields(capture.read(1).get(0)), "setData", "/enhanced", "1", "failed", "-103");
+ assertArrayEquals(unicode, zk.getData("/enhanced", false, null));
+ zk.getChildren("/", false);
+ zk.exists("/enhanced", false);
+ capture.read(0);
+ }
+
+ @Test
+ public void testEnhancedCreateVariants() throws Exception {
+ System.setProperty(AuditHelperTest.ENHANCED_ENABLE, "true");
+ System.setProperty("zookeeper.extendedTypesEnabled", "true");
+ ZooKeeper zk = createClient();
+ Stat stat = new Stat();
+ zk.create("/create2", null, ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT, stat);
+ assertWrite(fields(capture.read(1).get(0)), "create", "/create2", "0", "committed", "0");
+ zk.create("/container", new byte[2], ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.CONTAINER);
+ Map container = fields(capture.read(1).get(0));
+ assertWrite(container, "create", "/container", "2", "committed", "0");
+ assertEquals("container", container.get("znode_type"));
+ String path = zk.create("/ttl-", new byte[3], ZooDefs.Ids.OPEN_ACL_UNSAFE,
+ CreateMode.PERSISTENT_SEQUENTIAL_WITH_TTL, null, 60000);
+ Map ttl = fields(capture.read(1).get(0));
+ assertWrite(ttl, "create", path, "3", "committed", "0");
+ assertEquals("persistent_sequential_with_ttl", ttl.get("znode_type"));
+
+ zk.create("/ttl", new byte[0], ZooDefs.Ids.OPEN_ACL_UNSAFE,
+ CreateMode.PERSISTENT_WITH_TTL, null, 60000);
+ assertWrite(fields(capture.read(1).get(0)), "create", "/ttl", "0", "committed", "0");
+ try {
+ zk.create("/ttl", new byte[4], ZooDefs.Ids.OPEN_ACL_UNSAFE,
+ CreateMode.PERSISTENT_WITH_TTL, null, 60000);
+ fail("Duplicate TTL create must fail");
+ } catch (KeeperException e) {
+ assertEquals(Code.NODEEXISTS, e.code());
+ }
+ Map failedTtl = fields(capture.read(1).get(0));
+ assertWrite(failedTtl, "create", "/ttl", "4", "failed", "-110");
+ assertEquals("persistent_with_ttl", failedTtl.get("znode_type"));
+ }
+
+ @Test
+ public void testEnhancedMultiUsesReturnedPathsAndCompleteIndexes() throws Exception {
+ System.setProperty(AuditHelperTest.ENHANCED_ENABLE, "true");
+ System.setProperty("zookeeper.extendedTypesEnabled", "true");
+ ZooKeeper zk = createClient();
+ List results = zk.multi(Arrays.asList(
+ Op.create("/same", new byte[1], ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT),
+ Op.delete("/same", -1),
+ Op.check("/", -1),
+ Op.create("/same", new byte[2], ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.EPHEMERAL),
+ Op.create("/seq-", new byte[3], ZooDefs.Ids.OPEN_ACL_UNSAFE,
+ CreateMode.PERSISTENT_SEQUENTIAL_WITH_TTL, 60000)));
+ assertEquals(5, results.size());
+ String finalPath = ((OpResult.CreateResult) results.get(4)).getPath();
+ assertFalse("/seq-".equals(finalPath));
+ List logs = capture.read(4);
+ assertWrite(fields(logs.get(0)), "create", "/same", "1", "committed", "0");
+ assertEquals("persistent", fields(logs.get(0)).get("znode_type"));
+ assertWrite(fields(logs.get(1)), "delete", "/same", null, "committed", "0");
+ assertWrite(fields(logs.get(2)), "create", "/same", "2", "committed", "0");
+ assertEquals("ephemeral", fields(logs.get(2)).get("znode_type"));
+ assertEquals("3", fields(logs.get(2)).get("multi_index"));
+ assertWrite(fields(logs.get(3)), "create", finalPath, "3", "committed", "0");
+ assertEquals("4", fields(logs.get(3)).get("multi_index"));
+ assertEquals("persistent_sequential_with_ttl", fields(logs.get(3)).get("znode_type"));
+ assertArrayEquals(new byte[3], zk.getData(finalPath, false, null));
+ }
+
+ @Test
+ public void testEnhancedFailedMultiMatchesAtomicRollback() throws Exception {
+ System.setProperty(AuditHelperTest.ENHANCED_ENABLE, "true");
+ ZooKeeper zk = createClient();
+ try {
+ zk.multi(Arrays.asList(
+ Op.create("/rolled", new byte[1], ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT),
+ Op.check("/", -1),
+ Op.setData("/missing", new byte[3], -1),
+ Op.create("/later", new byte[2], ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT)));
+ fail("Missing node must abort the multi");
+ } catch (KeeperException e) {
+ assertEquals(Code.NONODE, e.code());
+ }
+ assertNull(zk.exists("/rolled", false));
+ assertNull(zk.exists("/later", false));
+ List logs = capture.read(4);
+ assertWrite(fields(logs.get(0)), "multiOperation", null, null, "failed", "-101");
+ assertWrite(fields(logs.get(1)), "create", "/rolled", "1", "rolled_back", "0");
+ assertWrite(fields(logs.get(2)), "setData", "/missing", "3", "failed", "-101");
+ assertWrite(fields(logs.get(3)), "create", "/later", "2", "rolled_back", "-2");
+ assertEquals("0", fields(logs.get(1)).get("multi_index"));
+ assertEquals("2", fields(logs.get(2)).get("multi_index"));
+ assertEquals("3", fields(logs.get(3)).get("multi_index"));
}
-}
+ @Test
+ public void testFailedCheckRollsBackMutationsWithoutACheckEvent() throws Exception {
+ System.setProperty(AuditHelperTest.ENHANCED_ENABLE, "true");
+ ZooKeeper zk = createClient();
+ try {
+ zk.multi(Arrays.asList(
+ Op.create("/rolled", new byte[1], ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT),
+ Op.check("/absent", -1),
+ Op.setData("/rolled", new byte[2], -1)));
+ fail("Missing check target must abort the multi");
+ } catch (KeeperException e) {
+ assertEquals(Code.NONODE, e.code());
+ }
+ assertNull(zk.exists("/rolled", false));
+ List logs = capture.read(3);
+ assertWrite(fields(logs.get(0)), "multiOperation", null, null, "failed", "-101");
+ assertWrite(fields(logs.get(1)), "create", "/rolled", "1", "rolled_back", "0");
+ assertWrite(fields(logs.get(2)), "setData", "/rolled", "2", "rolled_back", "-2");
+ assertEquals("2", fields(logs.get(2)).get("multi_index"));
+ }
+
+ @Test
+ public void testEnhancedAclRedactionDoesNotChangeStoredAcl() throws Exception {
+ System.setProperty(AuditHelperTest.ENHANCED_ENABLE, "true");
+ ZooKeeper zk = createClient();
+ String credentials = "alice:synthetic-password";
+ String digest = DigestAuthenticationProvider.generateDigest(credentials);
+ zk.addAuthInfo("digest", credentials.getBytes(StandardCharsets.UTF_8));
+ zk.create("/acl", new byte[0], ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT);
+ capture.read(1);
+ List acls = Collections.singletonList(new ACL(ZooDefs.Perms.ALL, new Id("digest", digest)));
+ zk.setACL("/acl", acls, -1);
+ String log = capture.read(1).get(0);
+ assertWrite(fields(log), "setAcl", "/acl", null, "committed", "0");
+ assertEquals("digest:alice:cdrwa", fields(log).get("acl"));
+ assertFalse(log.contains(digest));
+ assertFalse(log.contains("synthetic-password"));
+ assertEquals(acls, zk.getACL("/acl", new Stat()));
+ capture.read(0);
+ }
+
+ @Test
+ public void testRegisteredCustomUserIsRedactedOnlyInEnhancedMode() throws Exception {
+ String property = ProviderRegistry.AUTHPROVIDER_PROPERTY_PREFIX + "c1-audit-user";
+ String previous = System.getProperty(property);
+ System.setProperty(property, CredentialAuthenticationProvider.class.getName());
+ ProviderRegistry.initialize();
+ try {
+ System.setProperty(AuditHelperTest.ENHANCED_ENABLE, "true");
+ ZooKeeper zk = createClient();
+ zk.addAuthInfo("audit-test-custom", "alice:synthetic-password".getBytes(StandardCharsets.UTF_8));
+ zk.create("/custom-user", new byte[1], ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT);
+ String enhanced = capture.read(1).get(0);
+ assertWrite(fields(enhanced), "create", "/custom-user", "1", "committed", "0");
+ assertFalse(enhanced.contains("synthetic-password"));
+ List enhancedUsers = Arrays.asList(fields(enhanced).get("user").split(","));
+ Collections.sort(enhancedUsers);
+ assertEquals(Arrays.asList("127.0.0.1", "[redacted]"), enhancedUsers);
+
+ System.setProperty(AuditHelperTest.ENHANCED_ENABLE, "false");
+ zk.setData("/custom-user", new byte[2], -1);
+ Map legacy = fields(capture.read(1).get(0));
+ assertNull(legacy.get("schema_version"));
+ List legacyUsers = Arrays.asList(legacy.get("user").split(","));
+ Collections.sort(legacyUsers);
+ assertEquals(Arrays.asList("127.0.0.1", "alice:synthetic-password"), legacyUsers);
+ assertArrayEquals(new byte[2], zk.getData("/custom-user", false, null));
+ } finally {
+ ProviderRegistry.removeProvider("audit-test-custom");
+ AuditHelperTest.restoreProperty(property, previous);
+ }
+ }
+
+ @Test
+ public void testAuditFailureDoesNotRejectValidWrite() throws Exception {
+ System.setProperty(AuditHelperTest.ENHANCED_ENABLE, "true");
+ ZooKeeper zk = createClient();
+ AuditLogger failingLogger = event -> {
+ throw new IllegalStateException("synthetic audit sink failure");
+ };
+ Object previous = AuditHelperTest.replaceProviderField("auditLogger", failingLogger);
+ long before = AuditHelperTest.auditErrors();
+ try {
+ assertEquals("/valid", zk.create("/valid", new byte[2],
+ ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT));
+ assertArrayEquals(new byte[2], zk.getData("/valid", false, null));
+ assertEquals(before + 1, AuditHelperTest.auditErrors());
+ } finally {
+ AuditHelperTest.replaceProviderField("auditLogger", previous);
+ }
+ }
+
+ @Test
+ public void testAuditReporterFailureDoesNotRejectAppliedWrite() throws Exception {
+ ZooKeeper zk = createClient();
+ FailingCounter counter = new FailingCounter();
+ Counter previousCounter = AuditHelperTest.replaceAuditErrorCounter(counter);
+ Object previousLogger = AuditHelperTest.replaceProviderField("auditLogger", (AuditLogger) event -> {
+ throw new IllegalStateException("synthetic audit sink failure");
+ });
+ try {
+ for (String enhanced : Arrays.asList("false", "true")) {
+ System.setProperty(AuditHelperTest.ENHANCED_ENABLE, enhanced);
+ String path = "/reporter-" + enhanced;
+ Code replyError = null;
+ String created = null;
+ long before = counter.get();
+ try {
+ created = zk.create(path, new byte[2], ZooDefs.Ids.OPEN_ACL_UNSAFE, CreateMode.PERSISTENT);
+ } catch (KeeperException e) {
+ replyError = e.code();
+ }
+ assertArrayEquals(new byte[2], zk.getData(path, false, null));
+ assertNull("Audit reporting changed the reply after the write was applied", replyError);
+ assertEquals(path, created);
+ assertEquals("Do not retry a failing error reporter", before + 1, counter.get());
+ }
+ } finally {
+ AuditHelperTest.replaceProviderField("auditLogger", previousLogger);
+ AuditHelperTest.replaceAuditErrorCounter(previousCounter);
+ }
+ }
+}