|
12 | 12 | import com.launchdarkly.sdk.server.interfaces.DataSourceStatusProvider.Status; |
13 | 13 | import com.launchdarkly.sdk.server.interfaces.DataSourceStatusProvider.StatusListener; |
14 | 14 | import com.launchdarkly.sdk.server.subsystems.DataSourceUpdateSink; |
| 15 | +import com.launchdarkly.sdk.server.subsystems.DataSourceUpdateSinkV2; |
15 | 16 | import com.launchdarkly.sdk.server.subsystems.DataStore; |
| 17 | +import com.launchdarkly.sdk.server.subsystems.DataStoreTypes.ChangeSet; |
16 | 18 | import com.launchdarkly.sdk.server.subsystems.DataStoreTypes.DataKind; |
17 | 19 | import com.launchdarkly.sdk.server.subsystems.DataStoreTypes.FullDataSet; |
18 | 20 | import com.launchdarkly.sdk.server.subsystems.DataStoreTypes.ItemDescriptor; |
19 | 21 | import com.launchdarkly.sdk.server.subsystems.DataStoreTypes.KeyedItems; |
| 22 | +import com.launchdarkly.sdk.server.subsystems.TransactionalDataStore; |
20 | 23 | import com.launchdarkly.sdk.server.interfaces.DataStoreStatusProvider; |
21 | 24 | import com.launchdarkly.sdk.server.interfaces.FlagChangeEvent; |
22 | 25 | import com.launchdarkly.sdk.server.interfaces.FlagChangeListener; |
|
48 | 51 | * |
49 | 52 | * @since 4.11.0 |
50 | 53 | */ |
51 | | -final class DataSourceUpdatesImpl implements DataSourceUpdateSink { |
| 54 | +final class DataSourceUpdatesImpl implements DataSourceUpdateSink, DataSourceUpdateSinkV2 { |
52 | 55 | private final DataStore store; |
53 | 56 | private final EventBroadcasterImpl<FlagChangeListener, FlagChangeEvent> flagChangeEventNotifier; |
54 | 57 | private final EventBroadcasterImpl<StatusListener, Status> dataSourceStatusNotifier; |
@@ -365,4 +368,176 @@ private void onTimeout() { |
365 | 368 | private static String describeErrorCount(Map.Entry<ErrorInfo, Integer> entry) { |
366 | 369 | return entry.getKey() + " (" + entry.getValue() + (entry.getValue() == 1 ? " time" : " times") + ")"; |
367 | 370 | } |
| 371 | + |
| 372 | + @Override |
| 373 | + public boolean apply(ChangeSet<ItemDescriptor> changeSet) { |
| 374 | + if (store instanceof TransactionalDataStore) { |
| 375 | + return applyToTransactionalStore((TransactionalDataStore) store, changeSet); |
| 376 | + } |
| 377 | + |
| 378 | + // Legacy update path for non-transactional stores |
| 379 | + return applyToLegacyStore(changeSet); |
| 380 | + } |
| 381 | + |
| 382 | + private boolean applyToTransactionalStore(TransactionalDataStore transactionalDataStore, |
| 383 | + ChangeSet<ItemDescriptor> changeSet) { |
| 384 | + Map<DataKind, Map<String, ItemDescriptor>> oldData; |
| 385 | + // Getting the old values requires accessing the store, which can fail. |
| 386 | + // If there is a failure to read the store, then we stop treating it as a failure. |
| 387 | + try { |
| 388 | + oldData = getOldDataIfFlagChangeListeners(); |
| 389 | + } catch (RuntimeException e) { |
| 390 | + reportStoreFailure(e); |
| 391 | + return false; |
| 392 | + } |
| 393 | + |
| 394 | + ChangeSet<ItemDescriptor> sortedChangeSet = DataModelDependencies.sortChangeset(changeSet); |
| 395 | + |
| 396 | + try { |
| 397 | + transactionalDataStore.apply(sortedChangeSet); |
| 398 | + lastStoreUpdateFailed = false; |
| 399 | + } catch (RuntimeException e) { |
| 400 | + reportStoreFailure(e); |
| 401 | + return false; |
| 402 | + } |
| 403 | + |
| 404 | + // Calling Apply implies that the data source is now in a valid state. |
| 405 | + updateStatus(State.VALID, null); |
| 406 | + |
| 407 | + Set<KindAndKey> changes = updateDependencyTrackerForChangesetAndDetermineChanges(oldData, sortedChangeSet); |
| 408 | + |
| 409 | + // Now, if we previously queried the old data because someone is listening for flag change events, compare |
| 410 | + // the versions of all items and generate events for those (and any other items that depend on them) |
| 411 | + if (changes != null) { |
| 412 | + sendChangeEvents(changes); |
| 413 | + } |
| 414 | + |
| 415 | + return true; |
| 416 | + } |
| 417 | + |
| 418 | + private boolean applyToLegacyStore(ChangeSet<ItemDescriptor> sortedChangeSet) { |
| 419 | + switch (sortedChangeSet.getType()) { |
| 420 | + case Full: |
| 421 | + return applyFullChangeSetToLegacyStore(sortedChangeSet); |
| 422 | + case Partial: |
| 423 | + return applyPartialChangeSetToLegacyStore(sortedChangeSet); |
| 424 | + case None: |
| 425 | + default: |
| 426 | + return true; |
| 427 | + } |
| 428 | + } |
| 429 | + |
| 430 | + private boolean applyFullChangeSetToLegacyStore(ChangeSet<ItemDescriptor> unsortedChangeset) { |
| 431 | + // Convert ChangeSet to FullDataSet for legacy init path |
| 432 | + return init(new FullDataSet<>(unsortedChangeset.getData())); |
| 433 | + } |
| 434 | + |
| 435 | + private boolean applyPartialChangeSetToLegacyStore(ChangeSet<ItemDescriptor> changeSet) { |
| 436 | + // Sorting isn't strictly required here, as upsert behavior didn't traditionally have it, |
| 437 | + // but it also doesn't hurt, and there could be cases where it results in slightly |
| 438 | + // greater store consistency for persistent stores. |
| 439 | + ChangeSet<ItemDescriptor> sortedChangeset = DataModelDependencies.sortChangeset(changeSet); |
| 440 | + |
| 441 | + for (Map.Entry<DataKind, KeyedItems<ItemDescriptor>> kindItemsPair: sortedChangeset.getData()) { |
| 442 | + for (Map.Entry<String, ItemDescriptor> item: kindItemsPair.getValue().getItems()) { |
| 443 | + boolean applySuccess = upsert(kindItemsPair.getKey(), item.getKey(), item.getValue()); |
| 444 | + if (!applySuccess) { |
| 445 | + return false; |
| 446 | + } |
| 447 | + } |
| 448 | + } |
| 449 | + // The upsert will update the store status in the case of a store failure. |
| 450 | + // The application of the upserts does not set the store initialized. |
| 451 | + |
| 452 | + // Considering the store will be the same for the duration of the application |
| 453 | + // lifecycle we will not be applying a partial update to a store that didn't |
| 454 | + // already get a full update. The non-transactional store will also not support a selector. |
| 455 | + |
| 456 | + return true; |
| 457 | + } |
| 458 | + |
| 459 | + private Map<DataKind, Map<String, ItemDescriptor>> getOldDataIfFlagChangeListeners() { |
| 460 | + if (hasFlagChangeEventListeners()) { |
| 461 | + // Query the existing data if any, so that after the update we can send events for |
| 462 | + // whatever was changed |
| 463 | + Map<DataKind, Map<String, ItemDescriptor>> oldData = new HashMap<>(); |
| 464 | + for (DataKind kind: ALL_DATA_KINDS) { |
| 465 | + KeyedItems<ItemDescriptor> items = store.getAll(kind); |
| 466 | + oldData.put(kind, ImmutableMap.copyOf(items.getItems())); |
| 467 | + } |
| 468 | + return oldData; |
| 469 | + } else { |
| 470 | + return null; |
| 471 | + } |
| 472 | + } |
| 473 | + |
| 474 | + private Map<DataKind, Map<String, ItemDescriptor>> changeSetToMap(ChangeSet<ItemDescriptor> changeSet) { |
| 475 | + Map<DataKind, Map<String, ItemDescriptor>> ret = new HashMap<>(); |
| 476 | + for (Map.Entry<DataKind, KeyedItems<ItemDescriptor>> e: changeSet.getData()) { |
| 477 | + ret.put(e.getKey(), ImmutableMap.copyOf(e.getValue().getItems())); |
| 478 | + } |
| 479 | + return ret; |
| 480 | + } |
| 481 | + |
| 482 | + private Set<KindAndKey> updateDependencyTrackerForChangesetAndDetermineChanges( |
| 483 | + Map<DataKind, Map<String, ItemDescriptor>> oldDataMap, |
| 484 | + ChangeSet<ItemDescriptor> changeSet) { |
| 485 | + switch (changeSet.getType()) { |
| 486 | + case Full: |
| 487 | + return handleFullChangeset(oldDataMap, changeSet); |
| 488 | + case Partial: |
| 489 | + return handlePartialChangeset(oldDataMap, changeSet); |
| 490 | + case None: |
| 491 | + return null; |
| 492 | + default: |
| 493 | + return null; |
| 494 | + } |
| 495 | + } |
| 496 | + |
| 497 | + private Set<KindAndKey> handleFullChangeset( |
| 498 | + Map<DataKind, Map<String, ItemDescriptor>> oldDataMap, |
| 499 | + ChangeSet<ItemDescriptor> changeSet) { |
| 500 | + dependencyTracker.reset(); |
| 501 | + for (Map.Entry<DataKind, KeyedItems<ItemDescriptor>> kindEntry: changeSet.getData()) { |
| 502 | + DataKind kind = kindEntry.getKey(); |
| 503 | + for (Map.Entry<String, ItemDescriptor> itemEntry: kindEntry.getValue().getItems()) { |
| 504 | + String key = itemEntry.getKey(); |
| 505 | + dependencyTracker.updateDependenciesFrom(kind, key, itemEntry.getValue()); |
| 506 | + } |
| 507 | + } |
| 508 | + |
| 509 | + if (oldDataMap == null) { |
| 510 | + return null; |
| 511 | + } |
| 512 | + |
| 513 | + Map<DataKind, Map<String, ItemDescriptor>> newDataMap = changeSetToMap(changeSet); |
| 514 | + return computeChangedItemsForFullDataSet(oldDataMap, newDataMap); |
| 515 | + } |
| 516 | + |
| 517 | + private Set<KindAndKey> handlePartialChangeset( |
| 518 | + Map<DataKind, Map<String, ItemDescriptor>> oldDataMap, |
| 519 | + ChangeSet<ItemDescriptor> changeSet) { |
| 520 | + if (oldDataMap == null) { |
| 521 | + // Update dependencies but don't track changes when no listeners |
| 522 | + for (Map.Entry<DataKind, KeyedItems<ItemDescriptor>> kindEntry: changeSet.getData()) { |
| 523 | + DataKind kind = kindEntry.getKey(); |
| 524 | + for (Map.Entry<String, ItemDescriptor> itemEntry: kindEntry.getValue().getItems()) { |
| 525 | + dependencyTracker.updateDependenciesFrom(kind, itemEntry.getKey(), itemEntry.getValue()); |
| 526 | + } |
| 527 | + } |
| 528 | + return null; |
| 529 | + } |
| 530 | + |
| 531 | + Set<KindAndKey> affectedItems = new HashSet<>(); |
| 532 | + for (Map.Entry<DataKind, KeyedItems<ItemDescriptor>> kindEntry: changeSet.getData()) { |
| 533 | + DataKind kind = kindEntry.getKey(); |
| 534 | + for (Map.Entry<String, ItemDescriptor> itemEntry: kindEntry.getValue().getItems()) { |
| 535 | + String key = itemEntry.getKey(); |
| 536 | + dependencyTracker.updateDependenciesFrom(kind, key, itemEntry.getValue()); |
| 537 | + dependencyTracker.addAffectedItems(affectedItems, new KindAndKey(kind, key)); |
| 538 | + } |
| 539 | + } |
| 540 | + |
| 541 | + return affectedItems; |
| 542 | + } |
368 | 543 | } |
0 commit comments