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 @@ -5,6 +5,7 @@
import com.linkedin.openhouse.tables.readbridge.ColumnDefaultsSource;
import com.linkedin.openhouse.tables.readbridge.ReadBridgeConfigResolver;
import com.linkedin.openhouse.tables.readbridge.ReadBridgeStripProtection;
import com.linkedin.openhouse.tables.repository.OpenHouseInternalRepository;
import com.linkedin.openhouse.tables.toggle.TableFeatureToggle;
import org.springframework.beans.factory.ObjectProvider;
import org.springframework.context.annotation.Bean;
Expand All @@ -31,7 +32,9 @@ public ReadBridgeConfigResolver readBridgeConfigResolver(

@Bean
public ReadBridgeStripProtection readBridgeStripProtection(
ReadBridgeConfigResolver readBridgeConfigResolver) {
return new ReadBridgeStripProtection(readBridgeConfigResolver);
ReadBridgeConfigResolver readBridgeConfigResolver,
OpenHouseInternalRepository openHouseInternalRepository) {
return new ReadBridgeStripProtection(
readBridgeConfigResolver, openHouseInternalRepository::findSnapshotIds);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -7,18 +7,18 @@
import com.linkedin.openhouse.common.exception.DependencyUnavailableException;
import com.linkedin.openhouse.common.exception.TableConfigUnavailableException;
import com.linkedin.openhouse.tables.model.TableDto;
import com.linkedin.openhouse.tables.model.TableDtoPrimaryKey;
import java.util.ArrayList;
import java.util.Collections;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Objects;
import java.util.Set;
import java.util.function.Function;
import org.apache.iceberg.DataOperations;
import org.apache.iceberg.Snapshot;
import org.apache.iceberg.SnapshotParser;
import org.apache.iceberg.SnapshotRef;
import org.apache.iceberg.SnapshotRefParser;

/**
* Type 1 / Type 2 strip protection for column defaults. Default-aware clients send {@code
Expand All @@ -29,11 +29,10 @@
* table until OpenHouse rewrites are trusted or default-aware compaction exists. A lift/compaction
* flag is a follow-up, not this class.
*
* <p>The PUT includes the table's full snapshot list, so rewrite detection uses the main-branch
* snapshot, not a historical overwrite still in the list. A WAP / named-branch overwrite leaves
* {@code main} unchanged and is not gated. Do not infer the written ref by diffing ref maps — use
* commit deltas from #669 once they reach this path. See
* https://github.com/linkedin/openhouse/issues/693.
* <p>The PUT carries the table's full snapshot list, so Type 2 gates only the snapshots this commit
* adds: an overwrite/replace snapshot not already persisted on the table. That covers main, named
* branches, and staged WAP snapshots. A historical overwrite still in the list is not this commit,
* and a ref-only change (create branch or tag at an existing snapshot) adds nothing to gate.
*
* <p>Schema field objects are the JSON nodes that carry {@code id} — the same walk as the client
* overlay. On Iceberg schema JSON that is NestedField; {@code element-id} / {@code schema-id} are
Expand All @@ -52,9 +51,19 @@ private SchemaKeys() {}
}

private final ReadBridgeConfigResolver resolver;
private final Function<TableDtoPrimaryKey, Set<Long>> persistedSnapshotIds;

public ReadBridgeStripProtection(ReadBridgeConfigResolver resolver) {
/**
* @param persistedSnapshotIds every snapshot id in the table's persisted metadata. Consulted only
* when a stamped, ramped table's PUT carries an overwrite/replace snapshot. Its failures are
* not the request's fault and propagate as themselves.
*/
public ReadBridgeStripProtection(
ReadBridgeConfigResolver resolver,
Function<TableDtoPrimaryKey, Set<Long>> persistedSnapshotIds) {
this.resolver = Objects.requireNonNull(resolver, "resolver");
this.persistedSnapshotIds =
Objects.requireNonNull(persistedSnapshotIds, "persistedSnapshotIds");
}

/**
Expand All @@ -77,7 +86,7 @@ public TableDto prepare(TableDto existing, TableDto incoming) throws ColumnDefau
Map<Integer, String> incomingStamped = resolver.stampedColumnDefaults(incoming);
if (existing != null) {
rejectRemovedDefaults(previousStamped, incomingStamped, incoming);
rejectUnawareRewrite(previousStamped, incoming);
rejectUnawareRewrite(previousStamped, existing, incoming);
}
return stripInitialDefaults(incoming);
}
Expand Down Expand Up @@ -111,9 +120,10 @@ private void rejectRemovedDefaults(
* Type 2: overwrite/replace must send {@code initial-default} equal to the stamp. That handshake
* is trust, not proof the files were rewritten. Appends are not rewrites.
*/
private void rejectUnawareRewrite(Map<Integer, String> previousStamped, TableDto incoming)
private void rejectUnawareRewrite(
Map<Integer, String> previousStamped, TableDto existing, TableDto incoming)
throws ColumnDefaultException {
if (previousStamped.isEmpty() || !isRewrite(incoming)) {
if (previousStamped.isEmpty() || !isRewrite(existing, incoming)) {
return;
}
JsonNode schema = tree(incoming.getSchema(), incoming);
Expand Down Expand Up @@ -150,52 +160,48 @@ private TableDto stripInitialDefaults(TableDto incoming) throws ColumnDefaultExc
}

/**
* Main-branch snapshot only, so history in {@code jsonSnapshots} is not "this commit." WAP
* overwrite is therefore missed. Fix with #669 deltas, not a ref-map diff:
* https://github.com/linkedin/openhouse/issues/693
* A replace, or a commit that adds an overwrite/replace snapshot on any ref (or none, for a
* staged WAP snapshot). Snapshots already persisted on the table are history, not this commit.
* See https://github.com/linkedin/openhouse/issues/693.
*/
private boolean isRewrite(TableDto incoming) throws ColumnDefaultException {
private boolean isRewrite(TableDto existing, TableDto incoming) throws ColumnDefaultException {
if (incoming.isReplaceCommit() || incoming.isStageReplace()) {
return true;
}
List<String> jsonSnapshots = incoming.getJsonSnapshots();
if (jsonSnapshots == null || jsonSnapshots.isEmpty()) {
if (jsonSnapshots == null) {
return false;
}
String operation = currentSnapshot(incoming, jsonSnapshots).operation();
return DataOperations.OVERWRITE.equals(operation) || DataOperations.REPLACE.equals(operation);
}

private static Snapshot currentSnapshot(TableDto incoming, List<String> jsonSnapshots)
throws ColumnDefaultException {
Long mainId = mainSnapshotId(incoming);
if (mainId != null) {
for (String json : jsonSnapshots) {
Snapshot snapshot = snapshot(json, incoming);
if (snapshot.snapshotId() == mainId) {
return snapshot;
}
Set<Long> persisted = null;
for (String json : jsonSnapshots) {
Snapshot snapshot = snapshot(json, incoming);
String operation = snapshot.operation();
if (!DataOperations.OVERWRITE.equals(operation)
&& !DataOperations.REPLACE.equals(operation)) {
continue;
}
if (persisted == null) {
persisted = persistedSnapshotIds(existing);
}
if (!persisted.contains(snapshot.snapshotId())) {
return true;
}
throw ColumnDefaultException.unusable(
incoming, "main-branch snapshot is missing from the request", null);
}
return snapshot(jsonSnapshots.get(jsonSnapshots.size() - 1), incoming);
return false;
}

private static Long mainSnapshotId(TableDto incoming) throws ColumnDefaultException {
Map<String, String> snapshotRefs = incoming.getSnapshotRefs();
if (snapshotRefs == null) {
return null;
}
String main = snapshotRefs.get(SnapshotRef.MAIN_BRANCH);
if (main == null) {
return null;
}
try {
return SnapshotRefParser.fromJson(main).snapshotId();
} catch (RuntimeException e) {
throw ColumnDefaultException.unusable(incoming, "unreadable snapshot ref", e);
}
/**
* Read after {@code existing}. If another commit lands in between, this PUT's base version is
* stale and the commit is rejected regardless of this check.
*/
private Set<Long> persistedSnapshotIds(TableDto existing) {
return Objects.requireNonNull(
persistedSnapshotIds.apply(
TableDtoPrimaryKey.builder()
.databaseId(existing.getDatabaseId())
.tableId(existing.getTableId())
.build()),
"persisted snapshot ids");
}

private static Snapshot snapshot(String json, TableDto incoming) throws ColumnDefaultException {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@
import com.linkedin.openhouse.tables.model.TableDtoPrimaryKey;
import java.util.List;
import java.util.Optional;
import java.util.Set;
import org.springframework.data.domain.Page;
import org.springframework.data.domain.Pageable;
import org.springframework.data.repository.PagingAndSortingRepository;
Expand All @@ -28,6 +29,12 @@ public interface OpenHouseInternalRepository
*/
Optional<TableDto> findTableRefById(TableDtoPrimaryKey tableDtoPrimaryKey);

/**
* Ids of every snapshot in the table's current metadata, including snapshots no branch or tag
* references.
*/
Set<Long> findSnapshotIds(TableDtoPrimaryKey tableDtoPrimaryKey);

List<TableDtoPrimaryKey> findAllIds();

Page<TableDtoPrimaryKey> findAllIds(Pageable pageable);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -832,6 +832,17 @@ public Optional<TableDto> findTableRefById(TableDtoPrimaryKey tableDtoPrimaryKey
.build());
}

@Override
public Set<Long> findSnapshotIds(TableDtoPrimaryKey tableDtoPrimaryKey) {
Table table =
catalog.loadTable(
TableIdentifier.of(
tableDtoPrimaryKey.getDatabaseId(), tableDtoPrimaryKey.getTableId()));
Set<Long> snapshotIds = new HashSet<>();
table.snapshots().forEach(snapshot -> snapshotIds.add(snapshot.snapshotId()));
return snapshotIds;
}

// FIXME: Likely need a cache layer to avoid expensive tableScan.
@Timed(metricKey = MetricsConstant.REPO_TABLE_EXISTS_TIME)
@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@
import static org.hamcrest.Matchers.containsString;
import static org.hamcrest.Matchers.is;
import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.jsonPath;
import static org.springframework.test.web.servlet.result.MockMvcResultMatchers.status;

Expand All @@ -19,19 +20,32 @@
import com.linkedin.openhouse.common.test.cluster.PropertyOverrideContextInitializer;
import com.linkedin.openhouse.housetables.client.model.ToggleStatus;
import com.linkedin.openhouse.tables.api.spec.v0.request.CreateUpdateTableRequestBody;
import com.linkedin.openhouse.tables.api.spec.v0.request.IcebergSnapshotsRequestBody;
import com.linkedin.openhouse.tables.api.spec.v0.response.GetTableResponseBody;
import com.linkedin.openhouse.tables.mock.properties.AuthorizationPropertiesInitializer;
import com.linkedin.openhouse.tables.model.IcebergSnapshotsModelTestUtilities;
import com.linkedin.openhouse.tables.readbridge.ColumnDefaultException;
import com.linkedin.openhouse.tables.readbridge.ColumnDefaultsSource;
import com.linkedin.openhouse.tables.readbridge.ReadBridgeConfigResolver;
import com.linkedin.openhouse.tables.toggle.TableFeatureToggle;
import com.linkedin.openhouse.tables.toggle.model.TableToggleStatus;
import com.linkedin.openhouse.tables.toggle.repository.ToggleStatusesRepository;
import java.io.IOException;
import java.util.Arrays;
import java.util.Collections;
import java.util.HashMap;
import java.util.Map;
import java.util.UUID;
import java.util.stream.Collectors;
import org.apache.iceberg.DataFile;
import org.apache.iceberg.SchemaParser;
import org.apache.iceberg.Snapshot;
import org.apache.iceberg.SnapshotParser;
import org.apache.iceberg.SnapshotRef;
import org.apache.iceberg.SnapshotRefParser;
import org.apache.iceberg.Table;
import org.apache.iceberg.catalog.Catalog;
import org.apache.iceberg.catalog.TableIdentifier;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
Expand Down Expand Up @@ -91,6 +105,7 @@ ColumnDefaultsSource stubColumnDefaults() {
@Autowired private MockMvc mvc;
@Autowired private StorageManager storageManager;
@Autowired private ToggleStatusesRepository toggleStatusesRepository;
@Autowired private Catalog catalog;

private GetTableResponseBody created;
private TableToggleStatus toggleStatus;
Expand Down Expand Up @@ -212,15 +227,8 @@ public void putWithMatchingOverlay_doesNotPersistInitialDefault() throws Excepti

MvcResult get = getTable().andExpect(status().isOk()).andReturn();
GetTableResponseBody current = buildGetTableResponseBody(get);
ObjectMapper mapper = new ObjectMapper();
JsonNode root = mapper.readTree(current.getSchema());
for (JsonNode field : root.get("fields")) {
if (field.get("id").asInt() == 2) {
((ObjectNode) field).put("initial-default", "US");
}
}
GetTableResponseBody overlay =
current.toBuilder().schema(mapper.writeValueAsString(root)).build();
current.toBuilder().schema(withInitialDefault(current.getSchema())).build();

putTable(overlay).andExpect(status().isOk());

Expand Down Expand Up @@ -264,6 +272,104 @@ public void policyBypassAllowsDisablingCommittedOptIn() throws Exception {
.andExpect(jsonPath("$.config['" + CONFIG_KEY + "']").doesNotExist());
}

/**
* Type 2 gates the rewrite snapshots a snapshots PUT adds on any branch, not snapshots already on
* the table. https://github.com/linkedin/openhouse/issues/693
*/
@Test
public void snapshotsPut_gatesAddedBranchRewritesButNotRefsAtPersistedRewrites()
throws Exception {
created = create(uniqueTable("branch_rewrite"), Collections.singletonMap(ENABLED_PROP, "true"));
MvcResult current =
RequestAndValidateHelper.createTableAndValidateResponse(created, mvc, storageManager);
TableIdentifier identifier = TableIdentifier.of(created.getDatabaseId(), created.getTableId());

Table table = catalog.loadTable(identifier);
Snapshot append = table.newAppend().appendFile(dataFile(table)).apply();
current =
putSnapshots(current, false, branch("main", append), append)
.andExpect(status().isOk())
.andReturn();

table = catalog.loadTable(identifier);
Snapshot branchOverwrite =
table.newOverwrite().addFile(dataFile(table)).toBranch("feature").apply();
Map<String, String> branchWrite = branch("main", append);
branchWrite.putAll(branch("feature", branchOverwrite));
putSnapshots(current, false, branchWrite, append, branchOverwrite)
.andExpect(status().isBadRequest())
.andExpect(jsonPath("$.message", containsString("COLUMN_DEFAULT_REWRITE")));

Snapshot mainOverwrite = table.newOverwrite().addFile(dataFile(table)).apply();
current =
putSnapshots(current, true, branch("main", mainOverwrite), append, mainOverwrite)
.andExpect(status().isOk())
.andReturn();

Map<String, String> refOnly = branch("main", mainOverwrite);
refOnly.putAll(branch("feature", mainOverwrite));
refOnly.put(
"release",
SnapshotRefParser.toJson(SnapshotRef.tagBuilder(mainOverwrite.snapshotId()).build()));
putSnapshots(current, false, refOnly, append, mainOverwrite).andExpect(status().isOk());
assertTrue(catalog.loadTable(identifier).refs().get("release").isTag());
}

private ResultActions putSnapshots(
MvcResult current, boolean handshake, Map<String, String> refs, Snapshot... snapshots)
throws Exception {
CreateUpdateTableRequestBody envelope = buildCreateUpdateTableRequestBody(current);
if (handshake) {
envelope = envelope.toBuilder().schema(withInitialDefault(envelope.getSchema())).build();
}
IcebergSnapshotsRequestBody request =
IcebergSnapshotsRequestBody.builder()
.baseTableVersion(envelope.getBaseTableVersion())
.jsonSnapshots(
Arrays.stream(snapshots).map(SnapshotParser::toJson).collect(Collectors.toList()))
.snapshotRefs(refs)
.createUpdateTableRequestBody(envelope)
.build();
return mvc.perform(
MockMvcRequestBuilders.put(
String.format(
ValidationUtilities.CURRENT_MAJOR_VERSION_PREFIX
+ "/databases/%s/tables/%s/iceberg/v2/snapshots",
created.getDatabaseId(),
created.getTableId()))
.contentType(MediaType.APPLICATION_JSON)
.content(request.toJson())
.accept(MediaType.APPLICATION_JSON));
}

/** The default-aware client handshake: {@code initial-default} equal to the stamp. */
private static String withInitialDefault(String schemaJson) throws IOException {
ObjectMapper mapper = new ObjectMapper();
JsonNode root = mapper.readTree(schemaJson);
for (JsonNode field : root.get("fields")) {
if (field.get("id").asInt() == 2) {
((ObjectNode) field).put("initial-default", "US");
}
}
return mapper.writeValueAsString(root);
}

private DataFile dataFile(Table table) throws IOException {
return IcebergSnapshotsModelTestUtilities.createDummyDataFile(
storageManager.getDefaultStorage().getClient().getRootPrefix()
+ "/"
+ UUID.randomUUID()
+ ".orc",
table.spec());
}

private static Map<String, String> branch(String name, Snapshot snapshot) {
Map<String, String> refs = new HashMap<>();
refs.put(
name, SnapshotRefParser.toJson(SnapshotRef.branchBuilder(snapshot.snapshotId()).build()));
return refs;
}

private void activateHtsToggle(GetTableResponseBody table) {
toggleStatus =
TableToggleStatus.builder()
Expand Down
Loading