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