Skip to content
Closed
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 @@ -20,12 +20,17 @@
import com.couchbase.client.java.Cluster;
import com.couchbase.client.java.transactions.TransactionAttemptContext;
import io.flamingock.internal.common.core.audit.AuditPersistenceFactory;
import io.flamingock.internal.common.core.audit.AuditHistoryAppender;
import io.flamingock.internal.common.core.audit.JournalHistoryAppender;
import io.flamingock.internal.common.core.audit.AuditReader;
import io.flamingock.internal.common.core.audit.AuditEntry;
import io.flamingock.internal.common.core.journal.JournalEvent;
import io.flamingock.internal.common.core.context.ContextResolver;
import io.flamingock.internal.common.core.error.FlamingockException;
import io.flamingock.internal.common.core.feature.Features;
import io.flamingock.internal.core.configuration.community.CommunityConfigurable;
import io.flamingock.internal.core.external.store.CommunityAuditStore;
import io.flamingock.internal.core.external.store.HistoryAppenderProvider;
import io.flamingock.internal.core.context.BasicRuntimeContext;
import io.flamingock.internal.core.external.store.audit.community.CommunityAuditPersistence;
import io.flamingock.internal.core.external.store.lock.community.CommunityLockService;
import io.flamingock.internal.core.journal.JournalEventSequencer;
Expand All @@ -39,19 +44,19 @@
import io.flamingock.store.couchbase.internal.CouchbaseAuditPersistence;
import io.flamingock.store.couchbase.internal.CouchbaseAuditor;
import io.flamingock.store.couchbase.internal.CouchbaseJournalEventStore;
import io.flamingock.store.couchbase.internal.CouchbaseJournalWriter;
import io.flamingock.store.couchbase.internal.CouchbaseLockService;
import io.flamingock.externalsystem.couchbase.api.CouchbaseExternalSystem;

import java.util.Collections;
import java.util.Set;

public class CouchbaseAuditStore implements CommunityAuditStore {
public class CouchbaseAuditStore implements CommunityAuditStore, HistoryAppenderProvider {

private final CouchbaseExternalSystem targetSystem;
private final Cluster cluster;
private final String bucketName;
private RunnerId runnerId;
private CommunityConfigurable communityConfiguration;
private CouchbaseLockService lockService;
private Bucket bucket;
private String scopeName = CollectionIdentifier.DEFAULT_SCOPE;
Expand Down Expand Up @@ -122,42 +127,65 @@ public CouchbaseAuditStore withAutoCreate(boolean autoCreate) {
public void initialize(ContextResolver baseContext) {
this.validate();
runnerId = baseContext.getRequiredDependencyValue(RunnerId.class);
communityConfiguration = baseContext.getRequiredDependencyValue(CommunityConfigurable.class);

auditor = new CouchbaseAuditor(cluster, bucket);
journalEventStore = new CouchbaseJournalEventStore(cluster, bucket);
journalEventSequencerFactory = new JournalEventSequencerFactory(journalEventStore);

lockService = new CouchbaseLockService(cluster, bucket, TimeService.getDefault());
auditor.initialize(autoCreate, scopeName, auditRepositoryName);
lockService.initialize(autoCreate, scopeName, lockRepositoryName);
FeatureFlag.ifEnabled(Features.JOURNAL_EVENTS,
() -> journalEventStore.initialize(autoCreate, scopeName, journalRepositoryName));
}

@Override
public AuditPersistenceFactory<CommunityAuditPersistence> getPersistenceFactory() {
return stageId -> {
// Must run before forStream(stageId): forStream seeds the sequence from the last persisted event,
// which Couchbase reports as empty until the journal store is initialized. Idempotent and
// synchronized, so repeating it in CouchbaseAuditPersistence#doInitialize is safe.
FeatureFlag.ifEnabled(Features.JOURNAL_EVENTS, () -> journalEventStore.initialize(autoCreate, scopeName, journalRepositoryName));
JournalEventSequencer journalEventSequencer = journalEventSequencerFactory.forStream(stageId);
JournalEventSequencer journalEventSequencer = FeatureFlag.isEnabled(Features.JOURNAL_EVENTS, false)
? journalEventSequencerFactory.forStream(stageId) : null;
CouchbaseAuditPersistence persistence = new CouchbaseAuditPersistence(
communityConfiguration,
auditor,
journalEventStore,
journalEventSequencer,
targetSystem.getTxWrapper(),
scopeName,
auditRepositoryName,
journalRepositoryName,
autoCreate);
new CouchbaseJournalWriter(journalEventStore));
persistence.initialize(runnerId);
return persistence;
};
}

@Override
public AuditHistoryAppender getAuditHistoryAppender() {
return auditor::append;
}

@Override
public JournalHistoryAppender getJournalHistoryAppender() {
return (streamId, entry) -> {
if (!FeatureFlag.isEnabled(Features.JOURNAL_EVENTS, false)) {
throw new IllegalStateException("Journal events must be enabled to write journal history");
}
JournalEventSequencer sequencer = journalEventSequencerFactory.forStream(streamId);
CouchbaseJournalWriter writer = new CouchbaseJournalWriter(journalEventStore);
synchronized (sequencer) {
try {
JournalEvent<AuditEntry> event = sequencer.newEvent(entry);
io.flamingock.internal.util.Result result = targetSystem.getTxWrapper().wrapExecution(
new BasicRuntimeContext("write-journal-" + entry.getChangeId()), runtimeContext ->
writer.write(runtimeContext.getContext().getRequiredDependencyValue(
TransactionAttemptContext.class), event));
sequencer.confirm();
return result;
} catch (RuntimeException | Error failure) {
sequencer.markWriteOutcomeUncertain();
throw failure;
}
}
};
}

@Override
public AuditReader getAuditReader() {
auditor.initialize(autoCreate, scopeName, auditRepositoryName);
return () -> auditor.getAuditHistory();
}

Expand Down Expand Up @@ -198,19 +226,22 @@ private void validate() {
throw new FlamingockException("The 'lockRepositoryName' property is required.");
}

if (journalRepositoryName == null || journalRepositoryName.trim().isEmpty()) {
if (FeatureFlag.isEnabled(Features.JOURNAL_EVENTS, false)
&& (journalRepositoryName == null || journalRepositoryName.trim().isEmpty())) {
throw new FlamingockException("The 'journalRepositoryName' property is required.");
}

if (auditRepositoryName.trim().equalsIgnoreCase(lockRepositoryName.trim())) {
throw new FlamingockException("The 'auditRepositoryName' and 'lockRepositoryName' properties must not be the same.");
}

if (journalRepositoryName.trim().equalsIgnoreCase(auditRepositoryName.trim())) {
if (FeatureFlag.isEnabled(Features.JOURNAL_EVENTS, false)
&& journalRepositoryName.trim().equalsIgnoreCase(auditRepositoryName.trim())) {
throw new FlamingockException("The 'journalRepositoryName' and 'auditRepositoryName' properties must not be the same.");
}

if (journalRepositoryName.trim().equalsIgnoreCase(lockRepositoryName.trim())) {
if (FeatureFlag.isEnabled(Features.JOURNAL_EVENTS, false)
&& journalRepositoryName.trim().equalsIgnoreCase(lockRepositoryName.trim())) {
throw new FlamingockException("The 'journalRepositoryName' and 'lockRepositoryName' properties must not be the same.");
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -21,61 +21,32 @@
import io.flamingock.internal.common.core.feature.Features;
import io.flamingock.internal.common.core.journal.JournalEvent;
import io.flamingock.internal.common.core.external.ExecutionWrapper;
import io.flamingock.internal.core.configuration.community.CommunityConfigurable;
import io.flamingock.internal.core.context.BasicRuntimeContext;
import io.flamingock.internal.core.external.store.audit.community.AbstractCommunityAuditPersistence;
import io.flamingock.internal.core.journal.JournalEventSequencer;
import io.flamingock.internal.util.FeatureFlag;
import io.flamingock.internal.util.Result;
import io.flamingock.internal.util.id.RunnerId;

import java.util.List;

public class CouchbaseAuditPersistence extends AbstractCommunityAuditPersistence {

private final CouchbaseAuditor auditor;
private final CouchbaseJournalEventStore journalEventStore;
private final CouchbaseJournalWriter journalWriter;
private final JournalEventSequencer journalEventSequencer;
private final ExecutionWrapper txWrapper;
private final String scopeName;
private final String auditRepositoryName;
private final String journalRepositoryName;
private final boolean autoCreate;


public CouchbaseAuditPersistence(CommunityConfigurable localConfiguration,
CouchbaseAuditor auditor,
CouchbaseJournalEventStore journalEventStore,
public CouchbaseAuditPersistence(CouchbaseAuditor auditor,
JournalEventSequencer journalEventSequencer,
ExecutionWrapper txWrapper,
String scopeName,
String auditRepositoryName,
String journalRepositoryName,
boolean autoCreate) {
super(localConfiguration);
CouchbaseJournalWriter journalWriter) {
this.auditor = auditor;
this.journalEventStore = journalEventStore;
this.journalWriter = journalWriter;
this.journalEventSequencer = journalEventSequencer;
this.txWrapper = txWrapper;
this.scopeName = scopeName;
this.auditRepositoryName = auditRepositoryName;
this.journalRepositoryName = journalRepositoryName;
this.autoCreate = autoCreate;
}

@Override
protected void doInitialize(RunnerId runnerId) {
auditor.initialize(autoCreate, scopeName, auditRepositoryName);
// Creating the collection/indexes is what brings the journal collection into existence, so skipping
// this keeps it from ever appearing while the flag is off. It must stay in step with the append in
// writeEntry: skipping setup while still appending would let ctx.insert create the collection
// implicitly and without indexes, voiding the stream-position and eventId-lookup guarantees.
// Also repeated (idempotently) in CouchbaseAuditStore#getPersistenceFactory, which must run this
// before seeding the JournalEventSequencer via forStream(stageId).
FeatureFlag.ifEnabled(Features.JOURNAL_EVENTS, () -> journalEventStore.initialize(autoCreate, scopeName, journalRepositoryName));
}


@Override
public List<AuditEntry> getAuditHistory() {
return auditor.getAuditHistory();
Expand All @@ -88,24 +59,24 @@ public Result writeEntry(AuditEntry auditEntry) {
// without them, the audit record set is itself the history.
if (FeatureFlag.isEnabled(Features.JOURNAL_EVENTS)) {
RuntimeContext baseContext = new BasicRuntimeContext("write-changeState-" + auditEntry.getChangeId());
Result result = txWrapper.wrapExecution(baseContext, runtimeContext -> {
TransactionAttemptContext ctx = runtimeContext.getContext().getRequiredDependencyValue(TransactionAttemptContext.class);
JournalEvent<AuditEntry> journalEvent = journalEventSequencer.newEvent(auditEntry);
journalEventStore.contributeToTransaction(ctx, journalEvent);
return auditor.contributeToTransaction(ctx, auditEntry);
});
// Spends the stream position, and only a committed transaction attempt may reach this line. A
// normal return from wrapExecution does NOT in general mean commit — CouchbaseTxWrapper
// returns normally after a deliberate rollback too, when the operation's result is a FailedStep.
// It is sound here because this operation returns a Result, which can never be a FailedStep, so
// the only way to return normally is a committed attempt; a failing attempt is caught and
// rethrown as TransactionFailedException (see CouchbaseTxWrapper — it doesn't yet wrap that as
// DatabaseTransactionException, a known deviation from the ExecutionWrapper contract, tracked
// separately from this ticket). Keep that true: an operation that could return a failed step
// would silently burn a position and gap the stream, and a contiguous sequence is what lets a
// consumer tell "in flight" from "lost".
journalEventSequencer.confirm();
return result;
synchronized (journalEventSequencer) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Same issue as the DynamoDB version: possible null journalEventSequencer here if the flag flips on later, causing a crash. No null check like MongoDB has.

try {
JournalEvent<AuditEntry> event = journalEventSequencer.newEvent(auditEntry);
Result result = txWrapper.wrapExecution(baseContext, runtimeContext -> {
TransactionAttemptContext ctx = runtimeContext.getContext()
.getRequiredDependencyValue(TransactionAttemptContext.class);
journalWriter.write(ctx, event);
return auditor.contributeToTransaction(ctx, auditEntry);
});
// Result cannot be a FailedStep: a normal wrapper return means commit. Keep the
// transaction and confirmation under the same stream lock as journal-only writes.
journalEventSequencer.confirm();
return result;
} catch (RuntimeException | Error failure) {
journalEventSequencer.markWriteOutcomeUncertain();
throw failure;
}
}
} else {
return auditor.append(auditEntry);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -83,7 +83,7 @@ public synchronized void initialize(boolean autoCreate, String scopeName, String
* {@link #contributeToTransaction} would collapse them onto each other, discarding the very history being
* imported.
*/
Result append(AuditEntry auditEntry) {
public Result append(AuditEntry auditEntry) {

String key = toKey(auditEntry);
logger.debug("Saving audit entry with key {}", key);
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,36 @@
/*
* Copyright 2026 Flamingock (https://www.flamingock.io)
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package io.flamingock.store.couchbase.internal;

import com.couchbase.client.java.transactions.TransactionAttemptContext;
import io.flamingock.internal.common.core.audit.AuditEntry;
import io.flamingock.internal.common.core.journal.JournalEvent;
import io.flamingock.internal.util.Result;

/** Stages an event in the caller's transaction; confirmation belongs to the transaction owner. */
public class CouchbaseJournalWriter {

private final CouchbaseJournalEventStore store;

public CouchbaseJournalWriter(CouchbaseJournalEventStore store) {
this.store = store;
}

public Result write(TransactionAttemptContext context, JournalEvent<AuditEntry> event) {
store.contributeToTransaction(context, event);
return Result.OK();
}
}
Loading
Loading