diff --git a/hadoop-hdds/common/src/main/java/org/apache/hadoop/ozone/OzoneManagerVersion.java b/hadoop-hdds/common/src/main/java/org/apache/hadoop/ozone/OzoneManagerVersion.java index 9bd041c615c4..0038f1e93a21 100644 --- a/hadoop-hdds/common/src/main/java/org/apache/hadoop/ozone/OzoneManagerVersion.java +++ b/hadoop-hdds/common/src/main/java/org/apache/hadoop/ozone/OzoneManagerVersion.java @@ -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."); diff --git a/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/protocol/ClientProtocol.java b/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/protocol/ClientProtocol.java index cb553d33a9d3..41d1e3efebf6 100644 --- a/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/protocol/ClientProtocol.java +++ b/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/protocol/ClientProtocol.java @@ -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(); diff --git a/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/rpc/RpcClient.java b/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/rpc/RpcClient.java index 2eeb72c3e884..10210a7d2266 100644 --- a/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/rpc/RpcClient.java +++ b/hadoop-ozone/client/src/main/java/org/apache/hadoop/ozone/client/rpc/RpcClient.java @@ -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); } @@ -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); } @@ -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); } @@ -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 tags) throws IOException { if (omVersion.compareTo(OzoneManagerVersion.OBJECT_TAG) < 0) { @@ -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()); } @@ -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()); } @@ -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()); } @@ -2188,7 +2192,7 @@ private OpenKeySession newMultipartOpenKey( .setMultipartUploadPartNumber(partNumber) .setSortDatanodesInPipeline(sortDatanodesInPipeline) .setOwnerName(ownerName) - .setDerivedKeyPiggyBacking(derivedKeyPiggyBacking) + .setDerivedKeyPiggyBacking(shouldPiggybackDerivedKey(derivedKeyPiggyBacking)) .build(); return ozoneManagerClient.openKey(keyArgs); } diff --git a/hadoop-ozone/client/src/test/java/org/apache/hadoop/ozone/client/rpc/TestRpcClient.java b/hadoop-ozone/client/src/test/java/org/apache/hadoop/ozone/client/rpc/TestRpcClient.java index c7c4d343f6ce..d3d7328d5004 100644 --- a/hadoop-ozone/client/src/test/java/org/apache/hadoop/ozone/client/rpc/TestRpcClient.java +++ b/hadoop-ozone/client/src/test/java/org/apache/hadoop/ozone/client/rpc/TestRpcClient.java @@ -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; @@ -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; /** @@ -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 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 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 diff --git a/hadoop-ozone/s3gateway/src/main/java/org/apache/hadoop/ozone/s3/endpoint/EndpointBase.java b/hadoop-ozone/s3gateway/src/main/java/org/apache/hadoop/ozone/s3/endpoint/EndpointBase.java index 38aca6b6fe9c..17500fde290c 100644 --- a/hadoop-ozone/s3gateway/src/main/java/org/apache/hadoop/ozone/s3/endpoint/EndpointBase.java +++ b/hadoop-ozone/s3gateway/src/main/java/org/apache/hadoop/ozone/s3/endpoint/EndpointBase.java @@ -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; @@ -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); } @@ -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. * - *

In secure mode OM always returns the derived key for a signed upload, so + *

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. diff --git a/hadoop-ozone/s3gateway/src/main/java/org/apache/hadoop/ozone/s3/endpoint/ObjectEndpoint.java b/hadoop-ozone/s3gateway/src/main/java/org/apache/hadoop/ozone/s3/endpoint/ObjectEndpoint.java index e911b07eee16..4fa06f51b403 100644 --- a/hadoop-ozone/s3gateway/src/main/java/org/apache/hadoop/ozone/s3/endpoint/ObjectEndpoint.java +++ b/hadoop-ozone/s3gateway/src/main/java/org/apache/hadoop/ozone/s3/endpoint/ObjectEndpoint.java @@ -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); diff --git a/hadoop-ozone/s3gateway/src/test/java/org/apache/hadoop/ozone/s3/endpoint/TestEndpointBase.java b/hadoop-ozone/s3gateway/src/test/java/org/apache/hadoop/ozone/s3/endpoint/TestEndpointBase.java index 9865345a9162..0ca701a7c75c 100644 --- a/hadoop-ozone/s3gateway/src/test/java/org/apache/hadoop/ozone/s3/endpoint/TestEndpointBase.java +++ b/hadoop-ozone/s3gateway/src/test/java/org/apache/hadoop/ozone/s3/endpoint/TestEndpointBase.java @@ -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; @@ -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; /** @@ -209,4 +226,69 @@ private static Stream reservedInternalMetadataKeyPrefixCases() { RESERVED_USER_METADATA_KEY_PREFIX.toUpperCase(Locale.ROOT) + "cache-control"); } + static Stream 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); + } + } + } diff --git a/hadoop-ozone/s3gateway/src/test/java/org/apache/hadoop/ozone/s3/endpoint/TestObjectPut.java b/hadoop-ozone/s3gateway/src/test/java/org/apache/hadoop/ozone/s3/endpoint/TestObjectPut.java index 8e48c6ca1f92..655b6e2acacd 100644 --- a/hadoop-ozone/s3gateway/src/test/java/org/apache/hadoop/ozone/s3/endpoint/TestObjectPut.java +++ b/hadoop-ozone/s3gateway/src/test/java/org/apache/hadoop/ozone/s3/endpoint/TestObjectPut.java @@ -58,11 +58,16 @@ import static org.junit.jupiter.api.Assertions.assertThrows; import static org.mockito.Mockito.CALLS_REAL_METHODS; import static org.mockito.Mockito.any; +import static org.mockito.Mockito.anyBoolean; import static org.mockito.Mockito.anyInt; +import static org.mockito.Mockito.anyLong; +import static org.mockito.Mockito.anyMap; +import static org.mockito.Mockito.anyString; import static org.mockito.Mockito.doReturn; import static org.mockito.Mockito.doThrow; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.mockStatic; +import static org.mockito.Mockito.never; import static org.mockito.Mockito.spy; import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; @@ -93,12 +98,14 @@ import org.apache.hadoop.hdds.protocol.proto.HddsProtos; import org.apache.hadoop.ozone.OzoneConfigKeys; import org.apache.hadoop.ozone.OzoneConsts; +import org.apache.hadoop.ozone.OzoneManagerVersion; import org.apache.hadoop.ozone.client.BucketArgs; import org.apache.hadoop.ozone.client.OzoneBucket; import org.apache.hadoop.ozone.client.OzoneBucketStub; import org.apache.hadoop.ozone.client.OzoneClient; import org.apache.hadoop.ozone.client.OzoneKeyDetails; import org.apache.hadoop.ozone.client.OzoneVolume; +import org.apache.hadoop.ozone.client.protocol.ClientProtocol; import org.apache.hadoop.ozone.om.helpers.BucketLayout; import org.apache.hadoop.ozone.s3.HeaderPreprocessor; import org.apache.hadoop.ozone.s3.exception.OS3Exception; @@ -398,6 +405,36 @@ void testPutObjectRejectsMissingDerivedKeyInSecureMode() { assertKeyWasNotCommitted(); } + @ParameterizedTest + @ValueSource(booleans = {false, true}) + void testPutObjectRejectsUnsupportedOmBeforeOpeningKey(boolean datastream) throws Exception { + objectEndpoint.getOzoneConfiguration().setBoolean(OzoneConfigKeys.OZONE_SECURITY_ENABLED_KEY, true); + configureSignedChunkHeaders(CONTENT.length()); + doReturn(datastream).when(objectEndpoint).isDatastreamEnabled(); + doReturn(0L).when(objectEndpoint).getDatastreamMinLength(); + OzoneClient client = spy(objectEndpoint.getClient()); + ClientProtocol protocol = spy(client.getProxy()); + when(protocol.getOmVersion()).thenReturn(OzoneManagerVersion.GET_FILE_STATUS_REJECTS_OBS); + doReturn(protocol).when(client).getProxy(); + objectEndpoint.setClient(client); + // Rebuild the handler chain to use the spied endpoint and client. + objectEndpoint.init(); + byte[] payload = signedChunkedBody(CONTENT).getBytes(StandardCharsets.UTF_8); + when(headers.getHeaderString(HttpHeaders.CONTENT_LENGTH)).thenReturn(String.valueOf(payload.length)); + when(objectEndpoint.getContext().getMethod()).thenReturn(HttpMethod.PUT); + + try (MockedStatic streaming = mockStatic(ObjectEndpointStreaming.class); + ByteArrayInputStream body = new ByteArrayInputStream(payload)) { + assertErrorResponse(S3ErrorTable.NOT_IMPLEMENTED, () -> objectEndpoint.put(BUCKET_NAME, KEY_NAME, body)); + assertThat(body.available()).isEqualTo(payload.length); + verify(protocol).getOmVersion(); + verify(protocol, never()).createKey(anyString(), anyString(), anyString(), anyLong(), any(), anyMap(), anyMap(), + anyBoolean()); + streaming.verifyNoInteractions(); + assertKeyWasNotCommitted(); + } + } + @Test void testPutObjectRejectsInvalidSignedChunkTerminator() { configureSignedChunks(CONTENT.length()); diff --git a/hadoop-ozone/s3gateway/src/test/java/org/apache/hadoop/ozone/s3/endpoint/TestPartUpload.java b/hadoop-ozone/s3gateway/src/test/java/org/apache/hadoop/ozone/s3/endpoint/TestPartUpload.java index f8109676e148..5e3a625a8621 100644 --- a/hadoop-ozone/s3gateway/src/test/java/org/apache/hadoop/ozone/s3/endpoint/TestPartUpload.java +++ b/hadoop-ozone/s3gateway/src/test/java/org/apache/hadoop/ozone/s3/endpoint/TestPartUpload.java @@ -38,13 +38,20 @@ import static org.junit.jupiter.api.Assertions.assertNotNull; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.mockito.Mockito.CALLS_REAL_METHODS; +import static org.mockito.Mockito.anyBoolean; +import static org.mockito.Mockito.anyInt; +import static org.mockito.Mockito.anyLong; +import static org.mockito.Mockito.anyString; +import static org.mockito.Mockito.doReturn; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.mockStatic; +import static org.mockito.Mockito.never; import static org.mockito.Mockito.spy; import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; +import java.io.ByteArrayInputStream; import java.io.IOException; import java.io.InputStream; import java.nio.charset.StandardCharsets; @@ -59,10 +66,12 @@ import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.ozone.OzoneConfigKeys; import org.apache.hadoop.ozone.OzoneConsts; +import org.apache.hadoop.ozone.OzoneManagerVersion; import org.apache.hadoop.ozone.client.OzoneBucketStub; import org.apache.hadoop.ozone.client.OzoneClient; import org.apache.hadoop.ozone.client.OzoneClientStub; import org.apache.hadoop.ozone.client.OzoneMultipartUploadPartListParts; +import org.apache.hadoop.ozone.client.protocol.ClientProtocol; import org.apache.hadoop.ozone.s3.exception.OS3Exception; import org.apache.hadoop.ozone.s3.exception.S3ErrorTable; import org.apache.hadoop.ozone.s3.util.S3Consts; @@ -243,6 +252,47 @@ public void testPartUploadRejectsMissingDerivedKeyInSecureMode() throws Exceptio assertNoParts(uploadID, keyName); } + @ParameterizedTest + @ValueSource(strings = {STREAMING_AWS4_HMAC_SHA256_PAYLOAD, STREAMING_AWS4_HMAC_SHA256_PAYLOAD_TRAILER}) + void testPartUploadRejectsUnsupportedOmBeforeOpeningKey(String algorithm) throws Exception { + String keyName = UUID.randomUUID().toString(); + String uploadID = initiateMultipartUpload(rest, OzoneConsts.S3_BUCKET, keyName); + rest.getOzoneConfiguration().setBoolean(OzoneConfigKeys.OZONE_SECURITY_ENABLED_KEY, true); + configureSignedChunkHeaders(4); + when(headers.getHeaderString(X_AMZ_CONTENT_SHA256)).thenReturn(algorithm); + when(headers.getHeaderString(X_AMZ_TRAILER)).thenReturn("x-amz-checksum-crc32c"); + doReturn(0L).when(rest).getDatastreamMinLength(); + OzoneClient clientSpy = spy(client); + ClientProtocol protocol = spy(clientSpy.getProxy()); + when(protocol.getOmVersion()).thenReturn(OzoneManagerVersion.GET_FILE_STATUS_REJECTS_OBS); + doReturn(protocol).when(clientSpy).getProxy(); + rest.setClient(clientSpy); + // Rebuild the handler chain to use the spied endpoint and client. + rest.init(); + String chunkedBody = STREAMING_AWS4_HMAC_SHA256_PAYLOAD_TRAILER.equals(algorithm) + ? signedChunkedBodyWithTrailer("data", "x-amz-checksum-crc32c", "AAAAAA==") : signedChunkedBody("data"); + byte[] payload = chunkedBody.getBytes(StandardCharsets.UTF_8); + when(headers.getHeaderString(HttpHeaders.CONTENT_LENGTH)).thenReturn(String.valueOf(payload.length)); + when(rest.getContext().getMethod()).thenReturn(HttpMethod.PUT); + rest.queryParamsForTest().set(S3Consts.QueryParams.UPLOAD_ID, uploadID); + rest.queryParamsForTest().setInt(S3Consts.QueryParams.PART_NUMBER, 1); + + try (MockedStatic streaming = mockStatic(ObjectEndpointStreaming.class); + ByteArrayInputStream body = new ByteArrayInputStream(payload)) { + long failures = rest.getMetrics().getCreateMultipartKeyFailure(); + long successes = rest.getMetrics().getCreateMultipartKeySuccess(); + assertErrorResponse(S3ErrorTable.NOT_IMPLEMENTED, () -> rest.put(OzoneConsts.S3_BUCKET, keyName, body)); + assertThat(rest.getMetrics().getCreateMultipartKeyFailure()).isEqualTo(failures + 1); + assertThat(rest.getMetrics().getCreateMultipartKeySuccess()).isEqualTo(successes); + assertThat(body.available()).isEqualTo(payload.length); + verify(protocol).getOmVersion(); + verify(protocol, never()).createMultipartKey(anyString(), anyString(), anyString(), anyLong(), anyInt(), + anyString(), anyBoolean()); + streaming.verifyNoInteractions(); + assertNoParts(uploadID, keyName); + } + } + private void configureSignedChunks(int contentLength) throws IOException { OzoneBucketStub bucket = (OzoneBucketStub) client.getObjectStore() .getS3Bucket(OzoneConsts.S3_BUCKET);