From 735ec4a43834ea1e65aaa2c607230444b11b4441 Mon Sep 17 00:00:00 2001
From: Arkesh Mishra <118651144+arkmish@users.noreply.github.com>
Date: Fri, 25 Sep 2026 16:20:56 +0530
Subject: [PATCH 1/4] Add opt-in write size and outcome audit metadata (C1)
Decode typed requests without changing their buffers, cover TTL writes, correlate multi members by position and distinguish atomic rollback from commit. Add safe enhanced ACL metadata, deletion identity and observable audit failures while retaining legacy logging defaults.
Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
---
.../resources/markdown/zookeeperAuditLogs.md | 127 +++-
.../zookeeper/audit/AuditConstants.java | 2 +
.../apache/zookeeper/audit/AuditEvent.java | 25 +-
.../apache/zookeeper/audit/AuditHelper.java | 410 ++++++++----
.../zookeeper/audit/ZKAuditProvider.java | 63 +-
.../org/apache/zookeeper/server/DataTree.java | 15 +-
.../zookeeper/server/ServerMetrics.java | 2 +
.../zookeeper/audit/AuditEventTest.java | 53 ++
.../zookeeper/audit/AuditHelperTest.java | 597 ++++++++++++++++++
.../zookeeper/audit/Slf4JAuditLoggerTest.java | 121 ++--
.../audit/StandaloneServerAuditTest.java | 239 ++++++-
11 files changed, 1430 insertions(+), 224 deletions(-)
create mode 100644 zookeeper-server/src/test/java/org/apache/zookeeper/audit/AuditHelperTest.java
diff --git a/zookeeper-docs/src/main/resources/markdown/zookeeperAuditLogs.md b/zookeeper-docs/src/main/resources/markdown/zookeeperAuditLogs.md
index b3954ad676d..a58c8f19aa3 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,126 @@ 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.
+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. The
+enhanced gate is checked when emitting events; 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.
+
+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..c2811048aba 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,34 @@
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.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.server.util.AuthUtil;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -53,174 +64,295 @@ 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);
+ } 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);
+ 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;
+ }
}
- 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);
+ }
+ 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);
+ }
+ index++;
}
}
- private static void deserialize(Request request, Record record) throws IOException {
- request.request.rewind();
- ByteBufferInputStream.byteBuffer2Record(request.request.slice(), record);
- }
-
- 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 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 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 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 = AuthUtil.getUser(id);
+ }
}
+ value.append(id.getScheme()).append(':').append(user).append(':')
+ .append(ZKUtil.getPermString(acl.getPerms()));
}
- if (multiFailed) {
- log(request, rc.path, AuditConstants.OP_MULTI_OP, null,
- null, Result.FAILURE);
- }
+ return value.toString();
}
- 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) {
+ 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(request.getUsers(), operation, path, metadata.acl, metadata.createMode,
+ request.cnxn.getSessionIdHex(), request.cnxn.getHostAddress(), result,
+ metadata.dataLength, error, outcome, request.cxid, zxid, index);
}
+ private static void auditError(int type, Exception e) {
+ ServerMetrics.getMetrics().AUDIT_ERRORS.add(1);
+ LOG.error("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..701ab71a622 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,28 @@ 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;
+ }
+ logAuditEvent(createLogEvent(user, operation, znode, acl, createMode, session, ip, result,
+ dataLength, errorCode, outcome, cxid, zxid, multiIndex));
}
/**
@@ -78,6 +101,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);
return event;
}
@@ -86,6 +110,14 @@ 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) {
AuditEvent event = new AuditEvent(result);
event.addEntry(FieldName.SESSION, session);
event.addEntry(FieldName.USER, user);
@@ -94,9 +126,36 @@ 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);
return event;
}
+ private static void addMetadata(AuditEvent event, Integer dataLength, Integer errorCode, Outcome outcome,
+ Integer cxid, Long zxid, Integer multiIndex) {
+ if (isEnhancedAuditEnabled()) {
+ 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) {
+ ServerMetrics.getMetrics().AUDIT_ERRORS.add(1);
+ LOG.error("Failed to write audit log for operation {}", event.getValue(FieldName.OPERATION), e);
+ }
+ }
+
/**
* Add audit log for server start and register server stop log.
*/
@@ -119,7 +178,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..a85b1c9b9f7
--- /dev/null
+++ b/zookeeper-server/src/test/java/org/apache/zookeeper/audit/AuditHelperTest.java
@@ -0,0 +1,597 @@
+/*
+ * 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 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.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.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.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 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);
+ }
+ }
+
+ 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 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"));
+ }
+
+ 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..702857018d0 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,63 @@
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.data.ACL;
+import org.apache.zookeeper.data.Id;
+import org.apache.zookeeper.data.Stat;
+import org.apache.zookeeper.server.auth.DigestAuthenticationProvider;
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 +89,193 @@ 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());
+ }
+ 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());
}
- os.reset();
- return logs;
+ 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 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);
+ }
+ }
+}
From f37280b079d8bfb2342c814f5323bde40d64c237 Mon Sep 17 00:00:00 2001
From: Arkesh Mishra <118651144+arkmish@users.noreply.github.com>
Date: Fri, 25 Sep 2026 16:55:17 +0530
Subject: [PATCH 2/4] Redact untrusted enhanced audit users (C1 I1)
Use validated concrete built-in authentication providers for enhanced user extraction, redact custom and unknown identity representations, and preserve legacy behavior. Cover the registered default-getUserName leak with real client authentication.
Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
---
.../resources/markdown/zookeeperAuditLogs.md | 11 +++
.../apache/zookeeper/audit/AuditHelper.java | 59 +++++++++++++---
.../zookeeper/audit/AuditHelperTest.java | 69 +++++++++++++++++++
.../audit/StandaloneServerAuditTest.java | 34 +++++++++
4 files changed, 165 insertions(+), 8 deletions(-)
diff --git a/zookeeper-docs/src/main/resources/markdown/zookeeperAuditLogs.md b/zookeeper-docs/src/main/resources/markdown/zookeeperAuditLogs.md
index a58c8f19aa3..b9683a4fbb0 100644
--- a/zookeeper-docs/src/main/resources/markdown/zookeeperAuditLogs.md
+++ b/zookeeper-docs/src/main/resources/markdown/zookeeperAuditLogs.md
@@ -129,6 +129,17 @@ retained, valid built-in digest identities use the provider's username extractio
(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.
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 c2811048aba..855fd68f4e8 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
@@ -44,8 +44,10 @@
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.IPAuthenticationProvider;
import org.apache.zookeeper.server.auth.ProviderRegistry;
-import org.apache.zookeeper.server.util.AuthUtil;
+import org.apache.zookeeper.server.auth.SASLAuthenticationProvider;
+import org.apache.zookeeper.server.auth.X509AuthenticationProvider;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -54,6 +56,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);
@@ -95,7 +98,7 @@ public static void addAuditLog(Request request, ProcessTxnResult txnResult, bool
if (outcome != Outcome.COMMITTED || path == null) {
path = metadata.path == null ? path : metadata.path;
}
- log(request, txnResult, path, operation, metadata, result(outcome), error, outcome, null);
+ log(request, txnResult, path, operation, metadata, result(outcome), error, outcome, null, enhanced);
} catch (RuntimeException e) {
auditError(request.type, e);
}
@@ -159,7 +162,7 @@ private static void logMultiOperation(Request request, ProcessTxnResult rc, bool
// 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);
+ new RequestMetadata(), Result.FAILURE, error, Outcome.FAILED, null, enhanced);
if (!enhanced) {
return;
}
@@ -207,7 +210,7 @@ private static void logMultiOperation(Request request, ProcessTxnResult rc, bool
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);
+ subResult == null ? null : subResult.err, outcome, index, enhanced);
}
index++;
}
@@ -256,14 +259,14 @@ private static String safeAclToString(List acls) {
StringBuilder value = new StringBuilder();
for (ACL acl : acls) {
Id id = acl.getId();
- String user = "[redacted]";
+ 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 = AuthUtil.getUser(id);
+ user = provider.getUserName(id.getId());
}
}
value.append(id.getScheme()).append(':').append(user).append(':')
@@ -272,6 +275,45 @@ private static String safeAclToString(List acls) {
return value.toString();
}
+ 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;
+ }
+ 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 String createMode(int type, int flags) {
try {
return CreateMode.fromFlag(flags).name().toLowerCase(Locale.ROOT);
@@ -330,7 +372,8 @@ private static boolean isCreate(int type) {
}
private static void log(Request request, ProcessTxnResult rc, String path, String operation,
- RequestMetadata metadata, Result result, Integer error, Outcome outcome, Integer index) {
+ RequestMetadata metadata, Result result, Integer error, Outcome outcome,
+ Integer index, boolean enhanced) {
Long zxid = null;
if (request.getHdr() != null) {
zxid = request.getHdr().getZxid();
@@ -339,7 +382,7 @@ private static void log(Request request, ProcessTxnResult rc, String path, Strin
} else if (request.zxid >= 0) {
zxid = request.zxid;
}
- ZKAuditProvider.log(request.getUsers(), operation, path, metadata.acl, metadata.createMode,
+ 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);
}
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
index a85b1c9b9f7..7ddd49596f7 100644
--- a/zookeeper-server/src/test/java/org/apache/zookeeper/audit/AuditHelperTest.java
+++ b/zookeeper-server/src/test/java/org/apache/zookeeper/audit/AuditHelperTest.java
@@ -62,6 +62,8 @@
import org.apache.zookeeper.server.DataTree.ProcessTxnResult;
import org.apache.zookeeper.server.Request;
import org.apache.zookeeper.server.ServerCnxn;
+import org.apache.zookeeper.server.auth.AuthenticationProvider;
+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;
@@ -372,6 +374,45 @@ public void testMalformedDigestIdentityIsNotMistakenForAUser() throws Exception
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 testAuditDisabledSkipsRequestsAndProvider() throws Exception {
Object previous = replaceProviderField("auditEnabled", false);
@@ -539,6 +580,34 @@ static void assertWrite(Map fields, String operation, String pat
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 AuditCapture extends AppenderBase implements AutoCloseable {
private final Logger logger = (Logger) LoggerFactory.getLogger(Slf4jAuditLogger.class);
private final Level previousLevel = logger.getLevel();
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 702857018d0..3a0180f31a2 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
@@ -41,10 +41,12 @@
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.data.ACL;
import org.apache.zookeeper.data.Id;
import org.apache.zookeeper.data.Stat;
import org.apache.zookeeper.server.auth.DigestAuthenticationProvider;
+import org.apache.zookeeper.server.auth.ProviderRegistry;
import org.apache.zookeeper.test.ClientBase;
import org.junit.After;
import org.junit.AfterClass;
@@ -260,6 +262,38 @@ public void testEnhancedAclRedactionDoesNotChangeStoredAcl() throws Exception {
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");
From fd72bc7b4a7e179b264ac8a9d73972bf857b9189 Mon Sep 17 00:00:00 2001
From: Arkesh Mishra <118651144+arkmish@users.noreply.github.com>
Date: Fri, 25 Sep 2026 16:58:31 +0530
Subject: [PATCH 3/4] Isolate failures in audit error reporting (C1 I2)
Report audit failures through independently best-effort counter and diagnostic operations without recursive retries. Verify an applied write still receives success and multiple system deletions complete when both the audit sink and error counter fail.
Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
---
.../resources/markdown/zookeeperAuditLogs.md | 4 +
.../apache/zookeeper/audit/AuditHelper.java | 4 +-
.../zookeeper/audit/ZKAuditProvider.java | 14 +++-
.../zookeeper/audit/AuditHelperTest.java | 78 +++++++++++++++++++
.../audit/StandaloneServerAuditTest.java | 33 ++++++++
5 files changed, 129 insertions(+), 4 deletions(-)
diff --git a/zookeeper-docs/src/main/resources/markdown/zookeeperAuditLogs.md b/zookeeper-docs/src/main/resources/markdown/zookeeperAuditLogs.md
index b9683a4fbb0..8c0bd0e4234 100644
--- a/zookeeper-docs/src/main/resources/markdown/zookeeperAuditLogs.md
+++ b/zookeeper-docs/src/main/resources/markdown/zookeeperAuditLogs.md
@@ -221,6 +221,10 @@ 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
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 855fd68f4e8..6a1f5e13bc2 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
@@ -41,7 +41,6 @@
import org.apache.zookeeper.server.ByteBufferInputStream;
import org.apache.zookeeper.server.DataTree.ProcessTxnResult;
import org.apache.zookeeper.server.Request;
-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.IPAuthenticationProvider;
@@ -388,8 +387,7 @@ private static void log(Request request, ProcessTxnResult rc, String path, Strin
}
private static void auditError(int type, Exception e) {
- ServerMetrics.getMetrics().AUDIT_ERRORS.add(1);
- LOG.error("Failed to audit log request {}", type, e);
+ ZKAuditProvider.reportAuditError(LOG, "Failed to audit log request {}", type, e);
}
private static final class RequestMetadata {
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 701ab71a622..d108ffa65e9 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
@@ -151,8 +151,20 @@ 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);
- LOG.error("Failed to write audit log for operation {}", event.getValue(FieldName.OPERATION), e);
+ } 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.
}
}
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
index 7ddd49596f7..2f3e4fdd139 100644
--- a/zookeeper-server/src/test/java/org/apache/zookeeper/audit/AuditHelperTest.java
+++ b/zookeeper-server/src/test/java/org/apache/zookeeper/audit/AuditHelperTest.java
@@ -41,6 +41,7 @@
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;
@@ -53,6 +54,7 @@
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;
@@ -62,6 +64,7 @@
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.ProviderRegistry;
import org.apache.zookeeper.txn.CheckVersionTxn;
@@ -503,6 +506,57 @@ public void testLoggerFailureDoesNotInterruptSystemDeletion() throws Exception {
}
}
+ @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)));
}
@@ -549,6 +603,15 @@ static Object replaceProviderField(String name, Object value) throws ReflectiveO
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);
@@ -608,6 +671,21 @@ public boolean isValid(String id) {
}
}
+ 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();
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 3a0180f31a2..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
@@ -42,9 +42,11 @@
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;
@@ -312,4 +314,35 @@ public void testAuditFailureDoesNotRejectValidWrite() throws Exception {
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);
+ }
+ }
}
From 3dcb161260a5d98114a94e36e60fcaa012073164 Mon Sep 17 00:00:00 2001
From: Arkesh Mishra <118651144+arkmish@users.noreply.github.com>
Date: Fri, 25 Sep 2026 17:02:08 +0530
Subject: [PATCH 4/4] Keep one audit mode through metadata and emission (C1 I3)
Carry the captured enhanced mode into event construction without rereading the property. Preserve public delegating overloads and keep every multi parent/member on the same mode. Cover deterministic ACL off-to-on and multi mode transitions with subsequent genuine v2 emission.
Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com>
---
.../resources/markdown/zookeeperAuditLogs.md | 8 +-
.../apache/zookeeper/audit/AuditHelper.java | 2 +-
.../zookeeper/audit/ZKAuditProvider.java | 29 ++++-
.../zookeeper/audit/AuditHelperTest.java | 111 ++++++++++++++++++
4 files changed, 142 insertions(+), 8 deletions(-)
diff --git a/zookeeper-docs/src/main/resources/markdown/zookeeperAuditLogs.md b/zookeeper-docs/src/main/resources/markdown/zookeeperAuditLogs.md
index 8c0bd0e4234..4e03b13466d 100644
--- a/zookeeper-docs/src/main/resources/markdown/zookeeperAuditLogs.md
+++ b/zookeeper-docs/src/main/resources/markdown/zookeeperAuditLogs.md
@@ -210,8 +210,12 @@ 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. The
-enhanced gate is checked when emitting events; the base audit enablement retains
+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.
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 6a1f5e13bc2..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
@@ -383,7 +383,7 @@ private static void log(Request request, ProcessTxnResult rc, String path, Strin
}
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);
+ metadata.dataLength, error, outcome, request.cxid, zxid, index, enhanced);
}
private static void auditError(int type, Exception e) {
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 d108ffa65e9..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
@@ -90,8 +90,19 @@ public static void log(String user, String operation, String znode, String acl,
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));
+ dataLength, errorCode, outcome, cxid, zxid, multiIndex, enhanced));
}
/**
@@ -101,7 +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);
+ addMetadata(event, null, null, null, null, null, null, isEnhancedAuditEnabled());
return event;
}
@@ -118,6 +129,14 @@ static AuditEvent createLogEvent(String user, String operation, String znode, St
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);
@@ -126,13 +145,13 @@ 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);
+ 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) {
- if (isEnhancedAuditEnabled()) {
+ 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));
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
index 2f3e4fdd139..2dc2ab9f47a 100644
--- a/zookeeper-server/src/test/java/org/apache/zookeeper/audit/AuditHelperTest.java
+++ b/zookeeper-server/src/test/java/org/apache/zookeeper/audit/AuditHelperTest.java
@@ -66,6 +66,7 @@
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;
@@ -74,6 +75,7 @@
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;
@@ -416,6 +418,115 @@ public void testEnhancedUsersRedactUnknownAndMalformedIdentities() throws Except
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);