diff --git a/cluster/configs/src/main/java/com/linkedin/openhouse/cluster/configs/ClusterProperties.java b/cluster/configs/src/main/java/com/linkedin/openhouse/cluster/configs/ClusterProperties.java index 9d0c81876..ba894488b 100644 --- a/cluster/configs/src/main/java/com/linkedin/openhouse/cluster/configs/ClusterProperties.java +++ b/cluster/configs/src/main/java/com/linkedin/openhouse/cluster/configs/ClusterProperties.java @@ -96,4 +96,7 @@ public class ClusterProperties { // string @Value("${cluster.tables.allowed-client-name-values:}") private List allowedClientNameValues; + + @Value("${cluster.tables.iceberg-rest.enabled:false}") + private boolean clusterTablesIcebergRestEnabled; } diff --git a/docs/iceberg-rest-catalog.md b/docs/iceberg-rest-catalog.md new file mode 100644 index 000000000..e0b4ab641 --- /dev/null +++ b/docs/iceberg-rest-catalog.md @@ -0,0 +1,65 @@ +# Iceberg REST catalog + +OpenHouse exposes a read-only Apache Iceberg REST Catalog facade for new clients while preserving +the existing OpenHouse APIs and business behavior. + +## Enablement + +The facade is disabled by default. Enable it with: + +```properties +cluster.tables.iceberg-rest.enabled=true +``` + +`GET /v1/config` returns the `iceberg` route prefix and advertises only the implemented endpoints: + +- `GET /v1/{prefix}/namespaces/{namespace}/tables` +- `GET /v1/{prefix}/namespaces/{namespace}/tables/{table}` +- `HEAD /v1/{prefix}/namespaces/{namespace}/tables/{table}` + +## Architecture + +The OpenAPI-generated interfaces own the HTTP contract. `IcebergRestCatalogController` is a thin +Spring MVC adapter, and `IcebergRestApiHandler` translates the Iceberg protocol to existing +`TablesApiHandler` and `OpenHouseInternalCatalog` behavior. The facade does not add business rules +or change the existing OpenHouse endpoints. + +Iceberg response types use a narrowly scoped Spring `HttpMessageConverter`. Errors are translated +by controller-scoped advice into the standard Iceberg error envelope. + +## Compatibility and limitations + +- Only single-level namespaces are supported. +- The optional `warehouse` configuration hint does not select a different OpenHouse warehouse. +- List responses support opaque continuation tokens and page sizes from 1 through 1000. +- Table loads return all snapshots. The `snapshots=refs` projection is explicitly unsupported. +- The Iceberg 1.11 `referenced-by` query parameter is accepted and ignored. +- Access delegation may be requested, but this read-only version does not vend credentials. +- Conditional ETag responses are not currently emitted. +- Namespace, table-write, view, transaction, credential, and OAuth endpoints are not advertised. + +Existing OpenHouse APIs remain supported. Client migrations can therefore be incremental. + +## Observability and audit + +Spring Boot records the facade through the standard `http.server.requests` metrics, including URI, +status, and latency. Table reads delegate through `TablesApiHandler`, retaining existing +authorization, lock visibility, and table-read audit behavior. + +## Contract maintenance + +`spec/iceberg-rest-catalog-open-api.yaml` is the full Apache Iceberg REST OpenAPI. OpenHouse support +is opt-in: + +```yaml +operationId: listTables +x-openhouse-support: supported +``` + +Operations without the annotation are unsupported. The build codegens Spring interfaces only for +`supported` operations and generates `IcebergRestOpenHouseSupport.SUPPORTED_ENDPOINTS` for +`/v1/config`. Marking a new operation `supported` (or changing a supported signature) fails +compilation until the facade implements it. + +To upgrade Iceberg OpenAPI: merge the newer upstream YAML into the checked-in file, add +`x-openhouse-support: supported` where needed, then compile. diff --git a/infra/recipes/docker-compose/oh-only/docker-compose.yml b/infra/recipes/docker-compose/oh-only/docker-compose.yml index 823a21131..0ce12ecb4 100644 --- a/infra/recipes/docker-compose/oh-only/docker-compose.yml +++ b/infra/recipes/docker-compose/oh-only/docker-compose.yml @@ -4,6 +4,8 @@ services: extends: file: ../common/oh-services.yml service: openhouse-tables + environment: + CLUSTER_TABLES_ICEBERGREST_ENABLED: "true" volumes: - ./:/var/config/ depends_on: diff --git a/services/tables/src/main/java/com/linkedin/openhouse/tables/api/handler/IcebergRestApiHandler.java b/services/tables/src/main/java/com/linkedin/openhouse/tables/api/handler/IcebergRestApiHandler.java new file mode 100644 index 000000000..1a082ada7 --- /dev/null +++ b/services/tables/src/main/java/com/linkedin/openhouse/tables/api/handler/IcebergRestApiHandler.java @@ -0,0 +1,27 @@ +package com.linkedin.openhouse.tables.api.handler; + +import com.linkedin.openhouse.tables.generated.iceberg.model.CatalogConfig; +import com.linkedin.openhouse.tables.generated.iceberg.model.ListTablesResponse; +import org.apache.iceberg.rest.responses.LoadTableResponse; + +/** Protocol adapter between the generated Iceberg REST API and existing OpenHouse behavior. */ +public interface IcebergRestApiHandler { + + String ICEBERG_REST_PREFIX = "iceberg"; + + CatalogConfig getConfig(String warehouse); + + ListTablesResponse listTables( + String prefix, String namespace, String pageToken, Integer pageSize); + + LoadTableResponse loadTable( + String prefix, + String namespace, + String table, + String accessDelegation, + String ifNoneMatch, + String snapshots, + String referencedBy); + + void tableExists(String prefix, String namespace, String table); +} diff --git a/services/tables/src/main/java/com/linkedin/openhouse/tables/api/handler/impl/OpenHouseIcebergRestApiHandler.java b/services/tables/src/main/java/com/linkedin/openhouse/tables/api/handler/impl/OpenHouseIcebergRestApiHandler.java new file mode 100644 index 000000000..af1791dfc --- /dev/null +++ b/services/tables/src/main/java/com/linkedin/openhouse/tables/api/handler/impl/OpenHouseIcebergRestApiHandler.java @@ -0,0 +1,172 @@ +package com.linkedin.openhouse.tables.api.handler.impl; + +import static com.linkedin.openhouse.common.security.AuthenticationUtils.extractAuthenticatedUserPrincipal; + +import com.linkedin.openhouse.common.api.spec.ApiResponse; +import com.linkedin.openhouse.common.exception.NoSuchUserTableException; +import com.linkedin.openhouse.internal.catalog.OpenHouseInternalCatalog; +import com.linkedin.openhouse.tables.api.handler.IcebergRestApiHandler; +import com.linkedin.openhouse.tables.api.handler.TablesApiHandler; +import com.linkedin.openhouse.tables.api.spec.v0.response.GetAllTablesResponseBody; +import com.linkedin.openhouse.tables.api.spec.v0.response.GetTableResponseBody; +import com.linkedin.openhouse.tables.generated.iceberg.IcebergRestOpenHouseSupport; +import com.linkedin.openhouse.tables.generated.iceberg.model.CatalogConfig; +import com.linkedin.openhouse.tables.generated.iceberg.model.ListTablesResponse; +import java.nio.charset.StandardCharsets; +import java.util.Base64; +import java.util.Collections; +import java.util.LinkedHashSet; +import java.util.stream.Collectors; +import org.apache.iceberg.catalog.Namespace; +import org.apache.iceberg.catalog.TableIdentifier; +import org.apache.iceberg.exceptions.NoSuchNamespaceException; +import org.apache.iceberg.exceptions.NoSuchTableException; +import org.apache.iceberg.rest.CatalogHandlers; +import org.apache.iceberg.rest.RESTUtil; +import org.apache.iceberg.rest.responses.LoadTableResponse; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.data.domain.Page; +import org.springframework.stereotype.Component; + +/** Default Iceberg REST adapter backed by existing OpenHouse API handlers and catalog behavior. */ +@Component +@ConditionalOnProperty(value = "cluster.tables.iceberg-rest.enabled", havingValue = "true") +public class OpenHouseIcebergRestApiHandler implements IcebergRestApiHandler { + + static final int DEFAULT_PAGE_SIZE = 100; + static final int MAX_PAGE_SIZE = 1000; + private static final String PAGE_TOKEN_VERSION = "v1"; + + private final TablesApiHandler tablesApiHandler; + private final OpenHouseInternalCatalog openHouseInternalCatalog; + + public OpenHouseIcebergRestApiHandler( + TablesApiHandler tablesApiHandler, OpenHouseInternalCatalog openHouseInternalCatalog) { + this.tablesApiHandler = tablesApiHandler; + this.openHouseInternalCatalog = openHouseInternalCatalog; + } + + @Override + public CatalogConfig getConfig(String warehouse) { + return new CatalogConfig( + Collections.singletonMap("prefix", ICEBERG_REST_PREFIX), Collections.emptyMap()) + .endpoints(IcebergRestOpenHouseSupport.SUPPORTED_ENDPOINTS); + } + + @Override + public ListTablesResponse listTables( + String prefix, String namespace, String pageToken, Integer pageSize) { + validatePrefix(prefix); + Namespace icebergNamespace = decodeSingleLevelNamespace(namespace); + PageCursor cursor = decodePageToken(pageToken, pageSize); + ApiResponse response = + tablesApiHandler.searchTables( + icebergNamespace.level(0), cursor.getPage(), cursor.getPageSize(), "tableId"); + Page page = response.getResponseBody().getPageResults(); + LinkedHashSet identifiers = + page.getContent().stream() + .map(table -> TableIdentifier.of(icebergNamespace, table.getTableId())) + .collect(Collectors.toCollection(LinkedHashSet::new)); + String nextPageToken = + page.hasNext() ? encodePageToken(cursor.getPage() + 1, cursor.getPageSize()) : null; + return new ListTablesResponse().identifiers(identifiers).nextPageToken(nextPageToken); + } + + @Override + public LoadTableResponse loadTable( + String prefix, + String namespace, + String table, + String accessDelegation, + String ifNoneMatch, + String snapshots, + String referencedBy) { + validatePrefix(prefix); + if (snapshots != null && !"all".equals(snapshots)) { + throw new UnsupportedOperationException( + "The snapshots=refs projection is not supported by this catalog"); + } + // Iceberg 1.11 loadTable may send referenced-by for view-load chains; Phase 1 ignores it. + + Namespace icebergNamespace = decodeSingleLevelNamespace(namespace); + String databaseId = icebergNamespace.level(0); + try { + tablesApiHandler.getTable(databaseId, table, extractAuthenticatedUserPrincipal()); + } catch (NoSuchUserTableException e) { + throw new NoSuchTableException("Table does not exist: %s.%s", databaseId, table); + } + + return CatalogHandlers.loadTable( + openHouseInternalCatalog, TableIdentifier.of(icebergNamespace, table)); + } + + @Override + public void tableExists(String prefix, String namespace, String table) { + validatePrefix(prefix); + Namespace icebergNamespace = decodeSingleLevelNamespace(namespace); + String databaseId = icebergNamespace.level(0); + try { + tablesApiHandler.getTable(databaseId, table, extractAuthenticatedUserPrincipal()); + } catch (NoSuchUserTableException e) { + throw new NoSuchTableException("Table does not exist: %s.%s", databaseId, table); + } + } + + private static void validatePrefix(String prefix) { + if (!ICEBERG_REST_PREFIX.equals(prefix)) { + throw new IllegalArgumentException("Unsupported Iceberg REST prefix"); + } + } + + private static Namespace decodeSingleLevelNamespace(String encodedNamespace) { + Namespace namespace = RESTUtil.decodeNamespace(encodedNamespace); + if (namespace.isEmpty() || namespace.levels().length != 1) { + throw new NoSuchNamespaceException("Only single-level namespaces are supported"); + } + return namespace; + } + + private static PageCursor decodePageToken(String pageToken, Integer requestedPageSize) { + if (pageToken == null) { + return new PageCursor(0, validatePageSize(requestedPageSize)); + } + + try { + String decoded = new String(Base64.getUrlDecoder().decode(pageToken), StandardCharsets.UTF_8); + String[] parts = decoded.split(":", -1); + if (parts.length != 3 || !PAGE_TOKEN_VERSION.equals(parts[0])) { + throw new IllegalArgumentException("Invalid Iceberg REST page token"); + } + int page = Integer.parseInt(parts[1]); + int pageSize = validatePageSize(Integer.parseInt(parts[2])); + if (page < 1 || (requestedPageSize != null && requestedPageSize != pageSize)) { + throw new IllegalArgumentException("Invalid Iceberg REST page token"); + } + return new PageCursor(page, pageSize); + } catch (IllegalArgumentException e) { + throw new IllegalArgumentException("Invalid Iceberg REST page token", e); + } + } + + private static int validatePageSize(Integer requestedPageSize) { + int pageSize = requestedPageSize == null ? DEFAULT_PAGE_SIZE : requestedPageSize; + if (pageSize < 1 || pageSize > MAX_PAGE_SIZE) { + throw new IllegalArgumentException( + String.format("page-size must be between 1 and %s", MAX_PAGE_SIZE)); + } + return pageSize; + } + + private static String encodePageToken(int page, int pageSize) { + String value = String.format("%s:%s:%s", PAGE_TOKEN_VERSION, page, pageSize); + return Base64.getUrlEncoder() + .withoutPadding() + .encodeToString(value.getBytes(StandardCharsets.UTF_8)); + } + + @lombok.Value + private static class PageCursor { + int page; + int pageSize; + } +} diff --git a/services/tables/src/main/java/com/linkedin/openhouse/tables/controller/IcebergRestCatalogController.java b/services/tables/src/main/java/com/linkedin/openhouse/tables/controller/IcebergRestCatalogController.java new file mode 100644 index 000000000..06bbe107e --- /dev/null +++ b/services/tables/src/main/java/com/linkedin/openhouse/tables/controller/IcebergRestCatalogController.java @@ -0,0 +1,75 @@ +package com.linkedin.openhouse.tables.controller; + +import com.linkedin.openhouse.tables.api.handler.IcebergRestApiHandler; +import com.linkedin.openhouse.tables.generated.iceberg.api.CatalogApiApi; +import com.linkedin.openhouse.tables.generated.iceberg.api.ConfigurationApiApi; +import com.linkedin.openhouse.tables.generated.iceberg.model.CatalogConfig; +import com.linkedin.openhouse.tables.generated.iceberg.model.ListTablesResponse; +import io.swagger.v3.oas.annotations.Hidden; +import java.util.Optional; +import org.apache.iceberg.rest.responses.LoadTableResponse; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.http.ResponseEntity; +import org.springframework.web.bind.annotation.RestController; +import org.springframework.web.context.request.NativeWebRequest; + +/** + * Thin Spring MVC adapter for the generated read-only Iceberg REST contract. + * + *

Protocol translation and orchestration live in {@link IcebergRestApiHandler}; existing + * OpenHouse handlers and services remain the source of business behavior. + */ +@Hidden +@RestController +@ConditionalOnProperty(value = "cluster.tables.iceberg-rest.enabled", havingValue = "true") +public class IcebergRestCatalogController implements CatalogApiApi, ConfigurationApiApi { + + private final IcebergRestApiHandler icebergRestApiHandler; + + public IcebergRestCatalogController(IcebergRestApiHandler icebergRestApiHandler) { + this.icebergRestApiHandler = icebergRestApiHandler; + } + + @Override + public Optional getRequest() { + return Optional.empty(); + } + + @Override + public ResponseEntity getConfig(String warehouse) { + return ResponseEntity.ok(icebergRestApiHandler.getConfig(warehouse)); + } + + @Override + public ResponseEntity listTables( + String prefix, String namespace, String pageToken, Integer pageSize) { + return ResponseEntity.ok( + icebergRestApiHandler.listTables(prefix, namespace, pageToken, pageSize)); + } + + @Override + public ResponseEntity loadTable( + String prefix, + String namespace, + String table, + String xIcebergAccessDelegation, + String ifNoneMatch, + String snapshots, + String referencedBy) { + return ResponseEntity.ok( + icebergRestApiHandler.loadTable( + prefix, + namespace, + table, + xIcebergAccessDelegation, + ifNoneMatch, + snapshots, + referencedBy)); + } + + @Override + public ResponseEntity tableExists(String prefix, String namespace, String table) { + icebergRestApiHandler.tableExists(prefix, namespace, table); + return ResponseEntity.noContent().build(); + } +} diff --git a/services/tables/src/main/java/com/linkedin/openhouse/tables/controller/IcebergRestExceptionHandler.java b/services/tables/src/main/java/com/linkedin/openhouse/tables/controller/IcebergRestExceptionHandler.java new file mode 100644 index 000000000..103609bf5 --- /dev/null +++ b/services/tables/src/main/java/com/linkedin/openhouse/tables/controller/IcebergRestExceptionHandler.java @@ -0,0 +1,64 @@ +package com.linkedin.openhouse.tables.controller; + +import com.linkedin.openhouse.common.exception.RequestValidationFailureException; +import lombok.extern.slf4j.Slf4j; +import org.apache.iceberg.exceptions.ForbiddenException; +import org.apache.iceberg.exceptions.NoSuchNamespaceException; +import org.apache.iceberg.exceptions.NoSuchTableException; +import org.apache.iceberg.rest.responses.ErrorResponse; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.core.Ordered; +import org.springframework.core.annotation.Order; +import org.springframework.http.ResponseEntity; +import org.springframework.security.access.AccessDeniedException; +import org.springframework.web.bind.annotation.ExceptionHandler; +import org.springframework.web.bind.annotation.RestControllerAdvice; + +/** Scoped exception mapper for Iceberg REST endpoints. */ +@Order(Ordered.HIGHEST_PRECEDENCE) +@RestControllerAdvice(assignableTypes = IcebergRestCatalogController.class) +@ConditionalOnProperty(value = "cluster.tables.iceberg-rest.enabled", havingValue = "true") +@Slf4j +public class IcebergRestExceptionHandler { + + @ExceptionHandler(NoSuchTableException.class) + public ResponseEntity handleNoSuchTable(NoSuchTableException e) { + return errorResponse(404, e.getMessage(), NoSuchTableException.class.getSimpleName()); + } + + @ExceptionHandler(NoSuchNamespaceException.class) + public ResponseEntity handleNoSuchNamespace(NoSuchNamespaceException e) { + return errorResponse(404, e.getMessage(), NoSuchNamespaceException.class.getSimpleName()); + } + + @ExceptionHandler({RequestValidationFailureException.class, IllegalArgumentException.class}) + public ResponseEntity handleBadRequest(Exception e) { + return errorResponse(400, e.getMessage(), IllegalArgumentException.class.getSimpleName()); + } + + @ExceptionHandler(AccessDeniedException.class) + public ResponseEntity handleForbidden(AccessDeniedException e) { + return errorResponse(403, "Access denied", ForbiddenException.class.getSimpleName()); + } + + @ExceptionHandler(UnsupportedOperationException.class) + public ResponseEntity handleNotImplemented(UnsupportedOperationException e) { + return errorResponse(501, e.getMessage(), UnsupportedOperationException.class.getSimpleName()); + } + + @ExceptionHandler(Exception.class) + public ResponseEntity handleDefault(Exception e) { + log.error("Unhandled Iceberg REST request failure", e); + return errorResponse(500, "Internal server error", "InternalServerError"); + } + + private ResponseEntity errorResponse(int statusCode, String message, String type) { + ErrorResponse response = + ErrorResponse.builder() + .responseCode(statusCode) + .withMessage(message) + .withType(type) + .build(); + return ResponseEntity.status(statusCode).body(response); + } +} diff --git a/services/tables/src/main/java/com/linkedin/openhouse/tables/controller/IcebergRestHttpMessageConverter.java b/services/tables/src/main/java/com/linkedin/openhouse/tables/controller/IcebergRestHttpMessageConverter.java new file mode 100644 index 000000000..1ad538d14 --- /dev/null +++ b/services/tables/src/main/java/com/linkedin/openhouse/tables/controller/IcebergRestHttpMessageConverter.java @@ -0,0 +1,52 @@ +package com.linkedin.openhouse.tables.controller; + +import com.linkedin.openhouse.tables.generated.iceberg.model.CatalogConfig; +import com.linkedin.openhouse.tables.generated.iceberg.model.ListTablesResponse; +import java.io.IOException; +import java.io.OutputStream; +import java.nio.charset.StandardCharsets; +import org.apache.iceberg.rest.RESTResponse; +import org.springframework.http.HttpInputMessage; +import org.springframework.http.HttpOutputMessage; +import org.springframework.http.MediaType; +import org.springframework.http.converter.AbstractHttpMessageConverter; + +/** + * Spring {@link org.springframework.http.converter.HttpMessageConverter} for generated and runtime + * Iceberg response types. Uses {@link IcebergRestSerde} (Iceberg custom serializers, kebab-case) to + * write JSON because runtime Iceberg types do not follow JavaBean conventions. + * + *

This converter only handles writes (responses). Deserialization is not supported because we do + * not accept Iceberg REST request bodies through Spring MVC. + */ +public class IcebergRestHttpMessageConverter extends AbstractHttpMessageConverter { + + public IcebergRestHttpMessageConverter() { + super(MediaType.APPLICATION_JSON); + } + + @Override + protected boolean supports(Class clazz) { + return RESTResponse.class.isAssignableFrom(clazz) + || CatalogConfig.class.equals(clazz) + || ListTablesResponse.class.equals(clazz); + } + + @Override + public boolean canRead(Class clazz, MediaType mediaType) { + return false; + } + + @Override + protected Object readInternal(Class clazz, HttpInputMessage inputMessage) { + throw new UnsupportedOperationException("Iceberg REST request deserialization not supported"); + } + + @Override + protected void writeInternal(Object response, HttpOutputMessage outputMessage) + throws IOException { + OutputStream body = outputMessage.getBody(); + body.write(IcebergRestSerde.toJson(response).getBytes(StandardCharsets.UTF_8)); + body.flush(); + } +} diff --git a/services/tables/src/main/java/com/linkedin/openhouse/tables/controller/IcebergRestSerde.java b/services/tables/src/main/java/com/linkedin/openhouse/tables/controller/IcebergRestSerde.java new file mode 100644 index 000000000..e57b20df4 --- /dev/null +++ b/services/tables/src/main/java/com/linkedin/openhouse/tables/controller/IcebergRestSerde.java @@ -0,0 +1,33 @@ +package com.linkedin.openhouse.tables.controller; + +import com.fasterxml.jackson.annotation.JsonAutoDetect; +import com.fasterxml.jackson.annotation.PropertyAccessor; +import com.fasterxml.jackson.core.JsonFactory; +import com.fasterxml.jackson.core.JsonProcessingException; +import com.fasterxml.jackson.databind.DeserializationFeature; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.fasterxml.jackson.databind.PropertyNamingStrategies; +import org.apache.iceberg.rest.RESTSerializers; + +/** Serde helper for Iceberg REST payloads that require kebab-case and Iceberg serializers. */ +final class IcebergRestSerde { + + private static final ObjectMapper MAPPER = new ObjectMapper(new JsonFactory()); + + static { + MAPPER.setVisibility(PropertyAccessor.FIELD, JsonAutoDetect.Visibility.ANY); + MAPPER.configure(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES, false); + MAPPER.setPropertyNamingStrategy(new PropertyNamingStrategies.KebabCaseStrategy()); + RESTSerializers.registerAll(MAPPER); + } + + private IcebergRestSerde() {} + + static String toJson(Object payload) { + try { + return MAPPER.writeValueAsString(payload); + } catch (JsonProcessingException e) { + throw new IllegalStateException("Unable to serialize Iceberg REST response payload", e); + } + } +} diff --git a/services/tables/src/main/java/com/linkedin/openhouse/tables/controller/IcebergRestSerdeConfig.java b/services/tables/src/main/java/com/linkedin/openhouse/tables/controller/IcebergRestSerdeConfig.java new file mode 100644 index 000000000..a8295ac76 --- /dev/null +++ b/services/tables/src/main/java/com/linkedin/openhouse/tables/controller/IcebergRestSerdeConfig.java @@ -0,0 +1,21 @@ +package com.linkedin.openhouse.tables.controller; + +import java.util.List; +import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; +import org.springframework.context.annotation.Configuration; +import org.springframework.http.converter.HttpMessageConverter; +import org.springframework.web.servlet.config.annotation.WebMvcConfigurer; + +/** + * Registers the {@link IcebergRestHttpMessageConverter} so that Spring MVC can serialize typed + * Iceberg REST responses returned by {@link IcebergRestCatalogController}. + */ +@Configuration +@ConditionalOnProperty(value = "cluster.tables.iceberg-rest.enabled", havingValue = "true") +public class IcebergRestSerdeConfig implements WebMvcConfigurer { + + @Override + public void extendMessageConverters(List> converters) { + converters.add(0, new IcebergRestHttpMessageConverter()); + } +} diff --git a/services/tables/src/main/resources/application.properties b/services/tables/src/main/resources/application.properties index d6cfaa5d5..222be5d55 100644 --- a/services/tables/src/main/resources/application.properties +++ b/services/tables/src/main/resources/application.properties @@ -7,6 +7,7 @@ springdoc.swagger-ui.disable-swagger-default-url=true springdoc.swagger-ui.filter=true springdoc.swagger-ui.path=/tables/api-docs springdoc.swagger-ui.operationsSorter=method +cluster.tables.iceberg-rest.enabled=false server.tomcat.basedir=tomcat server.tomcat.accesslog.enabled=true server.tomcat.accesslog.rename-on-rotate=false diff --git a/services/tables/src/test/java/com/linkedin/openhouse/tables/mock/api/handler/impl/OpenHouseIcebergRestApiHandlerTest.java b/services/tables/src/test/java/com/linkedin/openhouse/tables/mock/api/handler/impl/OpenHouseIcebergRestApiHandlerTest.java new file mode 100644 index 000000000..047bdd52e --- /dev/null +++ b/services/tables/src/test/java/com/linkedin/openhouse/tables/mock/api/handler/impl/OpenHouseIcebergRestApiHandlerTest.java @@ -0,0 +1,143 @@ +package com.linkedin.openhouse.tables.mock.api.handler.impl; + +import static com.linkedin.openhouse.tables.api.handler.IcebergRestApiHandler.ICEBERG_REST_PREFIX; +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import com.linkedin.openhouse.common.api.spec.ApiResponse; +import com.linkedin.openhouse.internal.catalog.OpenHouseInternalCatalog; +import com.linkedin.openhouse.tables.api.handler.TablesApiHandler; +import com.linkedin.openhouse.tables.api.handler.impl.OpenHouseIcebergRestApiHandler; +import com.linkedin.openhouse.tables.api.spec.v0.response.GetAllTablesResponseBody; +import com.linkedin.openhouse.tables.api.spec.v0.response.GetTableResponseBody; +import com.linkedin.openhouse.tables.generated.iceberg.model.ListTablesResponse; +import java.util.Collections; +import org.apache.iceberg.BaseTable; +import org.apache.iceberg.PartitionSpec; +import org.apache.iceberg.Schema; +import org.apache.iceberg.SortOrder; +import org.apache.iceberg.TableMetadata; +import org.apache.iceberg.TableOperations; +import org.apache.iceberg.catalog.TableIdentifier; +import org.apache.iceberg.types.Types; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; +import org.springframework.data.domain.PageImpl; +import org.springframework.data.domain.PageRequest; +import org.springframework.http.HttpStatus; + +@ExtendWith(MockitoExtension.class) +public class OpenHouseIcebergRestApiHandlerTest { + + @Mock private TablesApiHandler tablesApiHandler; + @Mock private OpenHouseInternalCatalog openHouseInternalCatalog; + + private OpenHouseIcebergRestApiHandler handler; + + @BeforeEach + void setUp() { + handler = new OpenHouseIcebergRestApiHandler(tablesApiHandler, openHouseInternalCatalog); + } + + @Test + void configAdvertisesOnlyImplementedEndpoints() { + assertThat(handler.getConfig("openhouse").getOverrides()) + .containsEntry("prefix", ICEBERG_REST_PREFIX); + assertThat(handler.getConfig("openhouse").getEndpoints()) + .containsExactly( + "GET /v1/{prefix}/namespaces/{namespace}/tables", + "GET /v1/{prefix}/namespaces/{namespace}/tables/{table}", + "HEAD /v1/{prefix}/namespaces/{namespace}/tables/{table}"); + } + + @Test + void listTablesUsesOpaquePaginationToken() { + when(tablesApiHandler.searchTables("db", 0, 1, "tableId")) + .thenReturn(pageResponse("db", "t1", 0, 1, 2)); + + ListTablesResponse firstPage = handler.listTables(ICEBERG_REST_PREFIX, "db", null, 1); + + assertThat(firstPage.getIdentifiers()).containsExactly(TableIdentifier.of("db", "t1")); + assertThat(firstPage.getNextPageToken()).isNotBlank(); + + when(tablesApiHandler.searchTables("db", 1, 1, "tableId")) + .thenReturn(pageResponse("db", "t2", 1, 1, 2)); + ListTablesResponse secondPage = + handler.listTables(ICEBERG_REST_PREFIX, "db", firstPage.getNextPageToken(), null); + + assertThat(secondPage.getIdentifiers()).containsExactly(TableIdentifier.of("db", "t2")); + assertThat(secondPage.getNextPageToken()).isNull(); + } + + @Test + void rejectsInvalidPrefixAndPageToken() { + assertThatThrownBy(() -> handler.listTables("other", "db", null, null)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("prefix"); + assertThatThrownBy(() -> handler.listTables(ICEBERG_REST_PREFIX, "db", "invalid", null)) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("page token"); + } + + @Test + void rejectsUnsupportedSnapshotProjection() { + assertThatThrownBy( + () -> handler.loadTable(ICEBERG_REST_PREFIX, "db", "t1", null, null, "refs", null)) + .isInstanceOf(UnsupportedOperationException.class) + .hasMessageContaining("snapshots=refs"); + } + + @Test + void loadTableReusesExistingReadHandlerBeforeCatalogLoad() { + when(tablesApiHandler.getTable(eq("db"), eq("t1"), eq("undefined"))) + .thenReturn( + ApiResponse.builder() + .httpStatus(HttpStatus.OK) + .responseBody(GetTableResponseBody.builder().databaseId("db").tableId("t1").build()) + .build()); + TableMetadata metadata = testMetadata("hdfs://warehouse/db/t1"); + TableOperations operations = org.mockito.Mockito.mock(TableOperations.class); + when(operations.current()).thenReturn(metadata); + when(openHouseInternalCatalog.loadTable(TableIdentifier.of("db", "t1"))) + .thenReturn(new BaseTable(operations, "openhouse.db.t1")); + + assertThat( + handler + .loadTable(ICEBERG_REST_PREFIX, "db", "t1", null, null, "all", null) + .tableMetadata() + .location()) + .isEqualTo(metadata.location()); + verify(tablesApiHandler).getTable("db", "t1", "undefined"); + } + + private static ApiResponse pageResponse( + String databaseId, String tableId, int page, int size, int total) { + GetTableResponseBody table = + GetTableResponseBody.builder().databaseId(databaseId).tableId(tableId).build(); + return ApiResponse.builder() + .httpStatus(HttpStatus.OK) + .responseBody( + GetAllTablesResponseBody.builder() + .pageResults( + new PageImpl<>( + Collections.singletonList(table), PageRequest.of(page, size), total)) + .build()) + .build(); + } + + private static TableMetadata testMetadata(String location) { + Schema schema = new Schema(Types.NestedField.required(1, "id", Types.LongType.get())); + return TableMetadata.newTableMetadata( + schema, + PartitionSpec.unpartitioned(), + SortOrder.unsorted(), + location, + Collections.emptyMap()); + } +} diff --git a/services/tables/src/test/java/com/linkedin/openhouse/tables/mock/controller/IcebergRestCatalogControllerTest.java b/services/tables/src/test/java/com/linkedin/openhouse/tables/mock/controller/IcebergRestCatalogControllerTest.java new file mode 100644 index 000000000..af30fd547 --- /dev/null +++ b/services/tables/src/test/java/com/linkedin/openhouse/tables/mock/controller/IcebergRestCatalogControllerTest.java @@ -0,0 +1,175 @@ +package com.linkedin.openhouse.tables.mock.controller; + +import static com.linkedin.openhouse.tables.api.handler.IcebergRestApiHandler.ICEBERG_REST_PREFIX; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.ArgumentMatchers.nullable; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; +import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.jsonPath; +import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status; + +import com.linkedin.openhouse.tables.api.handler.IcebergRestApiHandler; +import com.linkedin.openhouse.tables.controller.IcebergRestCatalogController; +import com.linkedin.openhouse.tables.controller.IcebergRestExceptionHandler; +import com.linkedin.openhouse.tables.controller.IcebergRestHttpMessageConverter; +import com.linkedin.openhouse.tables.generated.iceberg.model.CatalogConfig; +import com.linkedin.openhouse.tables.generated.iceberg.model.ListTablesResponse; +import java.util.Arrays; +import java.util.Collections; +import java.util.LinkedHashSet; +import org.apache.iceberg.PartitionSpec; +import org.apache.iceberg.Schema; +import org.apache.iceberg.SortOrder; +import org.apache.iceberg.TableMetadata; +import org.apache.iceberg.catalog.TableIdentifier; +import org.apache.iceberg.exceptions.NoSuchTableException; +import org.apache.iceberg.rest.responses.LoadTableResponse; +import org.apache.iceberg.types.Types; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; +import org.springframework.http.converter.StringHttpMessageConverter; +import org.springframework.http.converter.json.MappingJackson2HttpMessageConverter; +import org.springframework.security.access.AccessDeniedException; +import org.springframework.test.web.servlet.MockMvc; +import org.springframework.test.web.servlet.request.MockMvcRequestBuilders; +import org.springframework.test.web.servlet.setup.MockMvcBuilders; + +@ExtendWith(MockitoExtension.class) +public class IcebergRestCatalogControllerTest { + + private MockMvc mvc; + + @Mock private IcebergRestApiHandler icebergRestApiHandler; + + @BeforeEach + public void setup() { + mvc = + MockMvcBuilders.standaloneSetup(new IcebergRestCatalogController(icebergRestApiHandler)) + .setControllerAdvice(new IcebergRestExceptionHandler()) + .setMessageConverters( + new IcebergRestHttpMessageConverter(), + new MappingJackson2HttpMessageConverter(), + new StringHttpMessageConverter()) + .build(); + } + + @Test + public void testConfigAdvertisesSupportedEndpoints() throws Exception { + when(icebergRestApiHandler.getConfig(nullable(String.class))) + .thenReturn( + new CatalogConfig( + Collections.singletonMap("prefix", ICEBERG_REST_PREFIX), Collections.emptyMap()) + .endpoints( + Collections.singletonList("GET /v1/{prefix}/namespaces/{namespace}/tables"))); + + mvc.perform(MockMvcRequestBuilders.get("/v1/config")) + .andExpect(status().isOk()) + .andExpect(jsonPath("$.overrides.prefix").value(ICEBERG_REST_PREFIX)) + .andExpect(jsonPath("$.endpoints[0]").exists()); + } + + @Test + public void testListTablesDelegatesTypedResponse() throws Exception { + when(icebergRestApiHandler.listTables( + eq(ICEBERG_REST_PREFIX), eq("db"), nullable(String.class), nullable(Integer.class))) + .thenReturn( + new ListTablesResponse() + .identifiers( + new LinkedHashSet<>( + Arrays.asList( + TableIdentifier.of("db", "tb1"), TableIdentifier.of("db", "tb2")))) + .nextPageToken("next")); + + mvc.perform( + MockMvcRequestBuilders.get("/v1/{prefix}/namespaces/db/tables", ICEBERG_REST_PREFIX) + .param("pageSize", "2")) + .andExpect(status().isOk()) + .andExpect(jsonPath("$.identifiers[0].namespace[0]").value("db")) + .andExpect(jsonPath("$.identifiers[1].name").value("tb2")) + .andExpect(jsonPath("$.next-page-token").value("next")); + } + + @Test + public void testLoadTableDelegatesTypedResponse() throws Exception { + TableMetadata metadata = testMetadata("hdfs://warehouse/db/tb1"); + when(icebergRestApiHandler.loadTable( + eq(ICEBERG_REST_PREFIX), + eq("db"), + eq("tb1"), + nullable(String.class), + nullable(String.class), + nullable(String.class), + nullable(String.class))) + .thenReturn(LoadTableResponse.builder().withTableMetadata(metadata).build()); + + mvc.perform( + MockMvcRequestBuilders.get( + "/v1/{prefix}/namespaces/db/tables/tb1", ICEBERG_REST_PREFIX)) + .andExpect(status().isOk()) + .andExpect(jsonPath("$.metadata-location").value(metadata.metadataFileLocation())) + .andExpect(jsonPath("$.metadata").exists()); + } + + @Test + public void testTypedNotFoundError() throws Exception { + when(icebergRestApiHandler.loadTable( + eq(ICEBERG_REST_PREFIX), + eq("db"), + eq("missing"), + nullable(String.class), + nullable(String.class), + nullable(String.class), + nullable(String.class))) + .thenThrow(new NoSuchTableException("Table does not exist")); + + mvc.perform( + MockMvcRequestBuilders.get( + "/v1/{prefix}/namespaces/db/tables/missing", ICEBERG_REST_PREFIX)) + .andExpect(status().isNotFound()) + .andExpect(jsonPath("$.error.code").value(404)) + .andExpect(jsonPath("$.error.type").value("NoSuchTableException")); + } + + @Test + public void testForbiddenErrorIsSanitized() throws Exception { + when(icebergRestApiHandler.loadTable( + eq(ICEBERG_REST_PREFIX), + eq("db"), + eq("private"), + nullable(String.class), + nullable(String.class), + nullable(String.class), + nullable(String.class))) + .thenThrow(new AccessDeniedException("sensitive policy details")); + + mvc.perform( + MockMvcRequestBuilders.get( + "/v1/{prefix}/namespaces/db/tables/private", ICEBERG_REST_PREFIX)) + .andExpect(status().isForbidden()) + .andExpect(jsonPath("$.error.message").value("Access denied")) + .andExpect(jsonPath("$.error.type").value("ForbiddenException")); + } + + @Test + public void testHeadDelegates() throws Exception { + mvc.perform( + MockMvcRequestBuilders.head( + "/v1/{prefix}/namespaces/db/tables/tb1", ICEBERG_REST_PREFIX)) + .andExpect(status().isNoContent()); + + verify(icebergRestApiHandler).tableExists(ICEBERG_REST_PREFIX, "db", "tb1"); + } + + private static TableMetadata testMetadata(String location) { + Schema schema = new Schema(Types.NestedField.required(1, "id", Types.LongType.get())); + return TableMetadata.newTableMetadata( + schema, + PartitionSpec.unpartitioned(), + SortOrder.unsorted(), + location, + Collections.emptyMap()); + } +} diff --git a/services/tables/src/test/java/com/linkedin/openhouse/tables/mock/controller/IcebergRestFeatureFlagTest.java b/services/tables/src/test/java/com/linkedin/openhouse/tables/mock/controller/IcebergRestFeatureFlagTest.java new file mode 100644 index 000000000..8a0de4815 --- /dev/null +++ b/services/tables/src/test/java/com/linkedin/openhouse/tables/mock/controller/IcebergRestFeatureFlagTest.java @@ -0,0 +1,41 @@ +package com.linkedin.openhouse.tables.mock.controller; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.Mockito.mock; + +import com.linkedin.openhouse.tables.api.handler.IcebergRestApiHandler; +import com.linkedin.openhouse.tables.controller.IcebergRestCatalogController; +import org.junit.jupiter.api.Test; +import org.springframework.boot.test.context.runner.ApplicationContextRunner; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.context.annotation.Import; + +public class IcebergRestFeatureFlagTest { + + private final ApplicationContextRunner contextRunner = + new ApplicationContextRunner().withUserConfiguration(TestConfiguration.class); + + @Test + void controllerIsDisabledByDefault() { + contextRunner.run( + context -> assertThat(context).doesNotHaveBean(IcebergRestCatalogController.class)); + } + + @Test + void controllerCanBeEnabled() { + contextRunner + .withPropertyValues("cluster.tables.iceberg-rest.enabled=true") + .run(context -> assertThat(context).hasSingleBean(IcebergRestCatalogController.class)); + } + + @Configuration(proxyBeanMethods = false) + @Import(IcebergRestCatalogController.class) + static class TestConfiguration { + + @Bean + IcebergRestApiHandler icebergRestApiHandler() { + return mock(IcebergRestApiHandler.class); + } + } +}