Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -67,6 +67,8 @@ public enum OzoneManagerVersion implements ComponentVersion {
+ "buckets server-side, so file system clients no longer need the "
+ "client-side InfoBucket layout check"),

S3_DERIVED_KEY(15, "OzoneManager version supporting derived signing keys for S3 chunk verification"),

FUTURE_VERSION(-1, "Used internally in the client when the server side is "
+ " newer and an unknown server version has arrived to the client.");

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1067,9 +1067,9 @@ TenantUserList listUsersInTenant(String tenantId, String prefix)
OzoneFsServerDefaults getServerDefaults() throws IOException;

/**
* Returns the negotiated Ozone Manager version for the connected cluster.
* In an HA cluster this is the minimum version across all OMs, so callers
* can safely gate client behavior on new server-side features.
* Returns the Ozone Manager version negotiated when the client was created.
* This is based on the minimum OM version advertised in the service list.
* The service list currently assumes all OM peers are at the same version.
* @return the effective Ozone Manager version.
*/
OzoneManagerVersion getOmVersion();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1431,7 +1431,7 @@ public OzoneOutputStream createKey(
OmKeyArgs.Builder builder = createWriteKeyArgsBuilder(volumeName,
bucketName, keyName, size, replicationConfig, metadata, tags);
builder.setOwnerName(ownerName);
builder.setDerivedKeyPiggyBacking(derivedKeyPiggyBacking);
builder.setDerivedKeyPiggyBacking(shouldPiggybackDerivedKey(derivedKeyPiggyBacking));
return openOutputStream(builder.build(), size);
}

Expand Down Expand Up @@ -1474,7 +1474,7 @@ public OzoneOutputStream createKeyIfNotExists(String volumeName,
bucketName, keyName, size, replicationConfig, metadata, tags);
builder.setExpectedDataGeneration(
OzoneConsts.EXPECTED_GEN_CREATE_IF_ABSENT);
builder.setDerivedKeyPiggyBacking(derivedKeyPiggyBacking);
builder.setDerivedKeyPiggyBacking(shouldPiggybackDerivedKey(derivedKeyPiggyBacking));
return openOutputStream(builder.build(), size);
}

Expand All @@ -1500,7 +1500,7 @@ public OzoneOutputStream rewriteKeyIfMatch(String volumeName,
OmKeyArgs.Builder builder = createWriteKeyArgsBuilder(volumeName,
bucketName, keyName, size, replicationConfig, metadata, tags);
builder.setExpectedETag(expectedETag);
builder.setDerivedKeyPiggyBacking(derivedKeyPiggyBacking);
builder.setDerivedKeyPiggyBacking(shouldPiggybackDerivedKey(derivedKeyPiggyBacking));
return openOutputStream(builder.build(), size);
}

Expand All @@ -1521,6 +1521,10 @@ private OzoneOutputStream openOutputStream(OmKeyArgs keyArgs, long size)
return createOutputStream(openKey);
}

private boolean shouldPiggybackDerivedKey(boolean requested) {
return requested && omVersion.compareTo(OzoneManagerVersion.S3_DERIVED_KEY) >= 0;
}

private void validateObjectTagsSupport(Map<String, String> tags)
throws IOException {
if (omVersion.compareTo(OzoneManagerVersion.OBJECT_TAG) < 0) {
Expand Down Expand Up @@ -1582,7 +1586,7 @@ public OzoneDataStreamOutput createStreamKey(
OmKeyArgs.Builder builder = createStreamKeyArgsBuilder(
volumeName, bucketName, keyName, size, replicationConfig, metadata,
tags);
builder.setDerivedKeyPiggyBacking(derivedKeyPiggyBacking);
builder.setDerivedKeyPiggyBacking(shouldPiggybackDerivedKey(derivedKeyPiggyBacking));
return openDataStreamOutput(builder.build());
}

Expand All @@ -1609,7 +1613,7 @@ public OzoneDataStreamOutput createStreamKeyIfNotExists(String volumeName,
tags);
builder.setExpectedDataGeneration(
OzoneConsts.EXPECTED_GEN_CREATE_IF_ABSENT);
builder.setDerivedKeyPiggyBacking(derivedKeyPiggyBacking);
builder.setDerivedKeyPiggyBacking(shouldPiggybackDerivedKey(derivedKeyPiggyBacking));
return openDataStreamOutput(builder.build());
}

Expand All @@ -1636,7 +1640,7 @@ public OzoneDataStreamOutput rewriteStreamKeyIfMatch(String volumeName,
volumeName, bucketName, keyName, size, replicationConfig, metadata,
tags);
builder.setExpectedETag(expectedETag);
builder.setDerivedKeyPiggyBacking(derivedKeyPiggyBacking);
builder.setDerivedKeyPiggyBacking(shouldPiggybackDerivedKey(derivedKeyPiggyBacking));
return openDataStreamOutput(builder.build());
}

Expand Down Expand Up @@ -2188,7 +2192,7 @@ private OpenKeySession newMultipartOpenKey(
.setMultipartUploadPartNumber(partNumber)
.setSortDatanodesInPipeline(sortDatanodesInPipeline)
.setOwnerName(ownerName)
.setDerivedKeyPiggyBacking(derivedKeyPiggyBacking)
.setDerivedKeyPiggyBacking(shouldPiggybackDerivedKey(derivedKeyPiggyBacking))
.build();
return ozoneManagerClient.openKey(keyArgs);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,8 +24,14 @@
import static org.junit.jupiter.api.Assertions.assertThrows;

import java.io.IOException;
import java.util.Arrays;
import java.util.Collections;
import java.util.LinkedList;
import java.util.List;
import java.util.concurrent.atomic.AtomicReference;
import java.util.stream.Stream;
import org.apache.hadoop.hdds.client.RatisReplicationConfig;
import org.apache.hadoop.hdds.client.ReplicationConfig;
import org.apache.hadoop.hdds.conf.OzoneConfiguration;
import org.apache.hadoop.hdds.protocol.proto.HddsProtos;
import org.apache.hadoop.hdds.scm.XceiverClientFactory;
Expand All @@ -37,12 +43,20 @@
import org.apache.hadoop.ozone.om.helpers.ServiceInfo;
import org.apache.hadoop.ozone.om.helpers.ServiceInfoEx;
import org.apache.hadoop.ozone.om.protocolPB.OmTransport;
import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos;
import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.CreateKeyRequest;
import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.OMRequest;
import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.OMResponse;
import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.ServiceListResponse;
import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.Type;
import org.apache.ozone.test.GenericTestUtils;
import org.apache.ozone.test.GenericTestUtils.LogCapturer;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.Arguments;
import org.junit.jupiter.params.provider.CsvSource;
import org.junit.jupiter.params.provider.EnumSource;
import org.junit.jupiter.params.provider.MethodSource;
import org.slf4j.event.Level;

/**
Expand Down Expand Up @@ -266,12 +280,104 @@ public void testListMultipartUploadsValidatesNames(String volume, String bucket,
}
}

private enum KeyWriteOperation {
CREATE, CREATE_IF_ABSENT, REWRITE_IF_MATCH,
STREAM, STREAM_IF_ABSENT, STREAM_IF_MATCH,
MULTIPART, MULTIPART_STREAM
}

static Stream<Arguments> derivedKeyPiggybackCases() {
return Arrays.stream(KeyWriteOperation.values()).flatMap(operation -> Stream.of(
Arguments.of(operation, OzoneManagerVersion.GET_FILE_STATUS_REJECTS_OBS.toProtoValue(), true, false),
Arguments.of(operation, OzoneManagerVersion.GET_FILE_STATUS_REJECTS_OBS.toProtoValue(), false, false),
Arguments.of(operation, OzoneManagerVersion.S3_DERIVED_KEY.toProtoValue(), true, true),
Arguments.of(operation, OzoneManagerVersion.S3_DERIVED_KEY.toProtoValue(), false, false),
Arguments.of(operation, OzoneManagerVersion.S3_DERIVED_KEY.toProtoValue() + 1, true, true),
Arguments.of(operation, OzoneManagerVersion.S3_DERIVED_KEY.toProtoValue() + 1, false, false)));
}

@ParameterizedTest
@MethodSource("derivedKeyPiggybackCases")
void testDerivedKeyPiggybackVersionGate(KeyWriteOperation operation, int omVersion,
boolean requested, boolean expected) throws IOException {
AtomicReference<CreateKeyRequest> captured = new AtomicReference<>();
OmTransport transport = new MockOmTransport() {
@Override
public OMResponse submitRequest(OMRequest request) throws IOException {
if (request.getCmdType() == Type.ServiceList) {
return super.submitRequest(request).toBuilder()
.setServiceListResponse(ServiceListResponse.newBuilder()
.addServiceInfo(OzoneManagerProtocolProtos.ServiceInfo.newBuilder()
.setNodeType(HddsProtos.NodeType.OM).setHostname("new-om")
.setOMVersion(OzoneManagerVersion.CURRENT.toProtoValue()))
.addServiceInfo(OzoneManagerProtocolProtos.ServiceInfo.newBuilder()
.setNodeType(HddsProtos.NodeType.OM).setHostname("other-om").setOMVersion(omVersion)))
.build();
}
if (request.getCmdType() == Type.CreateKey) {
captured.set(request.getCreateKeyRequest());
throw new IOException("Captured CreateKey request");
}
return super.submitRequest(request);
}
};
RpcClient client = createRpcClient(transport);
try {
ReplicationConfig replication = RatisReplicationConfig.getInstance(HddsProtos.ReplicationFactor.ONE);
IOException error = assertThrows(IOException.class, () -> {
switch (operation) {
case CREATE:
client.createKey("volume", "bucket", "key", 0, replication, Collections.emptyMap(), Collections.emptyMap(),
requested);
break;
case CREATE_IF_ABSENT:
client.createKeyIfNotExists("volume", "bucket", "key", 0, replication, Collections.emptyMap(),
Collections.emptyMap(), requested);
break;
case REWRITE_IF_MATCH:
client.rewriteKeyIfMatch("volume", "bucket", "key", 0, "etag", replication, Collections.emptyMap(),
Collections.emptyMap(), requested);
break;
case STREAM:
client.createStreamKey("volume", "bucket", "key", 0, replication, Collections.emptyMap(),
Collections.emptyMap(), requested);
break;
case STREAM_IF_ABSENT:
client.createStreamKeyIfNotExists("volume", "bucket", "key", 0, replication, Collections.emptyMap(),
Collections.emptyMap(), requested);
break;
case STREAM_IF_MATCH:
client.rewriteStreamKeyIfMatch("volume", "bucket", "key", 0, "etag", replication, Collections.emptyMap(),
Collections.emptyMap(), requested);
break;
case MULTIPART:
client.createMultipartKey("volume", "bucket", "key", 0, 1, "upload", requested);
break;
case MULTIPART_STREAM:
client.createMultipartStreamKey("volume", "bucket", "key", 0, 1, "upload", requested);
break;
default:
throw new IllegalArgumentException("Unexpected operation: " + operation);
}
});
assertThat(error).hasMessage("Captured CreateKey request");
assertThat(captured.get()).isNotNull();
assertThat(captured.get().getDerivedKeyPiggyBacking()).isEqualTo(expected);
} finally {
client.close();
}
}

private static RpcClient createRpcClient() throws IOException {
return createRpcClient(new MockOmTransport());
}

private static RpcClient createRpcClient(OmTransport transport) throws IOException {
OzoneConfiguration config = new OzoneConfiguration();
return new RpcClient(config, null) {
@Override
protected OmTransport createOmTransport(String omServiceId) {
return new MockOmTransport();
return transport;
}

@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -94,6 +94,7 @@
import org.apache.hadoop.hdds.conf.OzoneConfiguration;
import org.apache.hadoop.hdds.conf.StorageUnit;
import org.apache.hadoop.ozone.OzoneConsts;
import org.apache.hadoop.ozone.OzoneManagerVersion;
import org.apache.hadoop.ozone.OzoneSecurityUtil;
import org.apache.hadoop.ozone.audit.AuditAction;
import org.apache.hadoop.ozone.audit.AuditEventStatus;
Expand Down Expand Up @@ -878,6 +879,12 @@ protected S3ChunkInputStreamInfo getS3ChunkInputStreamInfo(
boolean verifyChunkSignature = signatureInfo.isSignPayload()
&& (STREAMING_AWS4_HMAC_SHA256_PAYLOAD.equals(amzContentSha256Header)
|| STREAMING_AWS4_HMAC_SHA256_PAYLOAD_TRAILER.equals(amzContentSha256Header));
if (verifyChunkSignature && OzoneSecurityUtil.isSecurityEnabled(getOzoneConfiguration())
&& getClientProtocol().getOmVersion().compareTo(OzoneManagerVersion.S3_DERIVED_KEY) < 0) {
OS3Exception ex = newError(S3ErrorTable.NOT_IMPLEMENTED, keyPath);
ex.setErrorMessage("The connected Ozone Manager does not support signed chunk verification.");
throw ex;
}
return new S3ChunkInputStreamInfo(multiDigestInputStream, effectiveLength, verifyChunkSignature);
}

Expand All @@ -886,7 +893,7 @@ protected S3ChunkInputStreamInfo getS3ChunkInputStreamInfo(
* the signing key OM derived (HDDS-15140). Called after the key is opened
* (when the derived key is available) and before the payload is read.
*
* <p>In secure mode OM always returns the derived key for a signed upload, so
* <p>In secure mode a supporting OM returns the derived key for a signed upload, so
* a missing key is treated as a server-side anomaly and the request is
* rejected rather than stored unverified. In non-secure mode there is no
* secret to verify against, so verification is skipped.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1103,7 +1103,7 @@ uploadID, getChunkSize(), multiDigestInputStream, perf, getHeaders(),
throw os3Exception;
}
throw newError(bucketName, key, ex);
} catch (IOException ex) {
} catch (IOException | RuntimeException ex) {
// Ensure we handle permission failures - these can surface as IOException wrapping OMException.
if (copyHeader != null) {
getMetrics().updateCopyObjectFailureStats(startNanos);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,10 +17,18 @@

package org.apache.hadoop.ozone.s3.endpoint;

import static org.apache.hadoop.ozone.OzoneConfigKeys.OZONE_SECURITY_ENABLED_KEY;
import static org.apache.hadoop.ozone.om.exceptions.OMException.ResultCodes;
import static org.apache.hadoop.ozone.s3.exception.S3ErrorTable.INVALID_ARGUMENT;
import static org.apache.hadoop.ozone.s3.signature.SignatureTestUtils.signatureInfo;
import static org.apache.hadoop.ozone.s3.util.S3Consts.CUSTOM_METADATA_HEADER_PREFIX;
import static org.apache.hadoop.ozone.s3.util.S3Consts.RESERVED_USER_METADATA_KEY_PREFIX;
import static org.apache.hadoop.ozone.s3.util.S3Consts.STREAMING_AWS4_HMAC_SHA256_PAYLOAD;
import static org.apache.hadoop.ozone.s3.util.S3Consts.STREAMING_AWS4_HMAC_SHA256_PAYLOAD_TRAILER;
import static org.apache.hadoop.ozone.s3.util.S3Consts.STREAMING_UNSIGNED_PAYLOAD_TRAILER;
import static org.apache.hadoop.ozone.s3.util.S3Consts.UNSIGNED_PAYLOAD;
import static org.apache.hadoop.ozone.s3.util.S3Consts.X_AMZ_CONTENT_SHA256;
import static org.apache.hadoop.ozone.s3.util.S3Consts.X_AMZ_TRAILER;
import static org.assertj.core.api.Assertions.assertThat;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
Expand All @@ -30,18 +38,27 @@
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.when;

import java.io.ByteArrayInputStream;
import java.nio.charset.StandardCharsets;
import java.util.Locale;
import java.util.Map;
import java.util.stream.Stream;
import javax.ws.rs.core.HttpHeaders;
import javax.ws.rs.core.MultivaluedHashMap;
import javax.ws.rs.core.MultivaluedMap;
import org.apache.hadoop.hdds.conf.OzoneConfiguration;
import org.apache.hadoop.ozone.OzoneConsts;
import org.apache.hadoop.ozone.OzoneManagerVersion;
import org.apache.hadoop.ozone.client.OzoneVolume;
import org.apache.hadoop.ozone.client.protocol.ClientProtocol;
import org.apache.hadoop.ozone.om.exceptions.OMException;
import org.apache.hadoop.ozone.s3.MultiDigestInputStream;
import org.apache.hadoop.ozone.s3.exception.OS3Exception;
import org.apache.hadoop.ozone.s3.exception.S3ErrorTable;
import org.apache.hadoop.ozone.s3.signature.SignatureTestUtils;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.Arguments;
import org.junit.jupiter.params.provider.MethodSource;

/**
Expand Down Expand Up @@ -209,4 +226,69 @@ private static Stream<String> reservedInternalMetadataKeyPrefixCases() {
RESERVED_USER_METADATA_KEY_PREFIX.toUpperCase(Locale.ROOT) + "cache-control");
}

static Stream<Arguments> signedChunkOmVersionCases() {
String payloadHash = SignatureTestUtils.sha256Hex(new byte[0], 0, 0);
return Stream.of(OzoneManagerVersion.DEFAULT_VERSION, OzoneManagerVersion.GET_FILE_STATUS_REJECTS_OBS,
OzoneManagerVersion.S3_DERIVED_KEY, OzoneManagerVersion.FUTURE_VERSION).flatMap(version -> {
boolean unsupported = version.compareTo(OzoneManagerVersion.S3_DERIVED_KEY) < 0;
return Stream.of(
Arguments.of(version, true, true, STREAMING_AWS4_HMAC_SHA256_PAYLOAD, unsupported, true),
Arguments.of(version, true, true, STREAMING_AWS4_HMAC_SHA256_PAYLOAD_TRAILER, unsupported, true),
Arguments.of(version, false, true, STREAMING_AWS4_HMAC_SHA256_PAYLOAD, false, true),
Arguments.of(version, false, true, STREAMING_AWS4_HMAC_SHA256_PAYLOAD_TRAILER, false, true),
Arguments.of(version, true, false, STREAMING_AWS4_HMAC_SHA256_PAYLOAD, false, false),
Arguments.of(version, true, false, STREAMING_AWS4_HMAC_SHA256_PAYLOAD_TRAILER, false, false),
Arguments.of(version, true, true, STREAMING_UNSIGNED_PAYLOAD_TRAILER, false, false),
Arguments.of(version, true, true, UNSIGNED_PAYLOAD, false, false),
Arguments.of(version, true, true, payloadHash, false, false));
});
}

@ParameterizedTest
@MethodSource("signedChunkOmVersionCases")
void testSignedChunksRequireSupportingOm(OzoneManagerVersion version, boolean secure, boolean signPayload,
String algorithm, boolean unsupported, boolean verificationRequired) throws Exception {
ClientProtocol protocol = mock(ClientProtocol.class);
when(protocol.getOmVersion()).thenReturn(version);
EndpointBase endpoint = new EndpointBase() {
@Override
protected ClientProtocol getClientProtocol() {
return protocol;
}
};
OzoneConfiguration conf = new OzoneConfiguration();
conf.setBoolean(OZONE_SECURITY_ENABLED_KEY, secure);
endpoint.setOzoneConfiguration(conf);
endpoint.setSignatureInfo(signatureInfo(signPayload));
HttpHeaders headers = mock(HttpHeaders.class);
when(headers.getHeaderString(X_AMZ_CONTENT_SHA256)).thenReturn(algorithm);
when(headers.getHeaderString(X_AMZ_TRAILER)).thenReturn("x-amz-checksum-crc32c");
endpoint.setHeaders(headers);

byte[] payload = "body".getBytes(StandardCharsets.UTF_8);
try (ByteArrayInputStream body = new ByteArrayInputStream(payload)) {
if (unsupported) {
OS3Exception ex = assertThrows(OS3Exception.class,
() -> endpoint.getS3ChunkInputStreamInfo(body, payload.length, "4", "key"));
assertThat(ex.getCode()).isEqualTo(S3ErrorTable.NOT_IMPLEMENTED.getCode());
assertThat(ex.getHttpCode()).isEqualTo(501);
assertThat(ex.getErrorMessage()).contains("Ozone Manager", "signed chunk verification");
} else {
EndpointBase.S3ChunkInputStreamInfo info = endpoint.getS3ChunkInputStreamInfo(body, payload.length, "4", "key");
try (MultiDigestInputStream stream = info.getMultiDigestInputStream()) {
assertThat(info.isChunkSignatureVerificationRequired()).isEqualTo(verificationRequired);
assertThat(info.getEffectiveLength()).isEqualTo(payload.length);
if (secure && verificationRequired) {
OS3Exception ex = assertThrows(OS3Exception.class, () -> endpoint.attachChunkValidator(info, "key", null));
assertThat(ex.getCode()).isEqualTo(S3ErrorTable.INTERNAL_ERROR.getCode());
} else {
endpoint.attachChunkValidator(info, "key", null);
}
assertThat(body.available()).isEqualTo(payload.length);
}
}
assertThat(body.available()).isEqualTo(payload.length);
}
}

}
Loading
Loading