Skip to content
Merged
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 @@ -3,8 +3,10 @@


import org.springframework.data.annotation.Id;
import org.springframework.data.mongodb.core.index.CompoundIndex;
import org.springframework.data.mongodb.core.mapping.Document;

@Document(collection = "HighLevelEvents")
@CompoundIndex(name = "projectId_timeStamp_idx", def = "{'projectId': 1, 'timeStamp': 1}")
public record HighLevelBusinessEvent(@Id String eventId,String projectId, int timeStamp, org.bson.Document projectEvent) {
}
Original file line number Diff line number Diff line change
Expand Up @@ -2,11 +2,13 @@

import com.fasterxml.jackson.databind.ObjectMapper;
import edu.stanford.protege.webprotege.common.ProjectEvent;
import edu.stanford.protege.webprotege.ipc.EventDispatcher;
import edu.stanford.protege.webprotegeeventshistory.dto.*;
import edu.stanford.protege.webprotegeeventshistory.sequence.SequenceService;
import edu.stanford.protege.webprotegeeventshistory.sequence.ProjectSequenceService;
import org.bson.Document;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.context.annotation.Lazy;
import org.springframework.stereotype.Service;
import org.springframework.transaction.annotation.Transactional;

Expand All @@ -24,27 +26,55 @@ public class HighLevelBusinessEventsService {

private final ObjectMapper objectMapper;

private final SequenceService sequenceService;
private final ProjectSequenceService projectSequenceService;

private final EventDispatcher eventDispatcher;

public HighLevelBusinessEventsService(HighLevelBusinessEventsRepository repository, ObjectMapper objectMapper, SequenceService sequenceService) {

// EventDispatcher is injected lazily to break a bean-creation cycle: the ipc EventDispatcher
// depends on RabbitMQEventsConfiguration, which eagerly wires every EventHandler, one of which
// (RegisterHighLevelBusinessEventHandler) depends back on this service.
public HighLevelBusinessEventsService(HighLevelBusinessEventsRepository repository,
ObjectMapper objectMapper,
ProjectSequenceService projectSequenceService,
@Lazy EventDispatcher eventDispatcher) {
this.repository = repository;
this.objectMapper = objectMapper;
this.sequenceService = sequenceService;
this.projectSequenceService = projectSequenceService;
this.eventDispatcher = eventDispatcher;
}


@Transactional
void registerEvent(PackagedProjectChangeEvent projectEvent) {
int seq;
try {
var nextDocument = objectMapper.convertValue(projectEvent, Document.class);

HighLevelBusinessEvent event = new HighLevelBusinessEvent(projectEvent.eventId().id(), projectEvent.projectId().id(), getTagFromNow(), nextDocument);
seq = projectSequenceService.next(projectEvent.projectId().id());
HighLevelBusinessEvent event = new HighLevelBusinessEvent(projectEvent.eventId().id(), projectEvent.projectId().id(), seq, nextDocument);
LOGGER.info("Logging event " + event);
repository.save(event);
} catch (Exception e) {
LOGGER.error("An error occurred when trying to save events", e);
return;
}
publishSequencedEvent(projectEvent, seq);
}

/**
* Re-publishes the persisted change bundle enriched with its sequence ordinal so the gateway
* can push it with truthful event tags. This runs after the archive write, so anything pushed
* is already durably fetchable. Unlike the persistence step above, a failure here is not
* swallowed: reliable delivery of this publish is tracked by #299, so it is allowed to surface
* rather than pass silently.
*/
private void publishSequencedEvent(PackagedProjectChangeEvent projectEvent, int seq) {
var sequencedEvent = new SequencedPackagedProjectChangeEvent(
projectEvent.projectId(),
projectEvent.eventId(),
seq,
projectEvent.projectEvents());
eventDispatcher.dispatchEvent(sequencedEvent);
}

ProjectEventsQueryResponse fetchEvents(ProjectEventsQueryRequest request) {
Expand All @@ -71,8 +101,4 @@ ProjectEventsQueryResponse fetchEvents(ProjectEventsQueryRequest request) {
response.events = new EventList<>(first, eventList, last);
return response;
}

private int getTagFromNow() {
return sequenceService.getNextHighLevelEventSequence();
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,45 @@
package edu.stanford.protege.webprotegeeventshistory.dto;

import com.fasterxml.jackson.annotation.JsonTypeName;
import edu.stanford.protege.webprotege.common.EventId;
import edu.stanford.protege.webprotege.common.ProjectEvent;
import edu.stanford.protege.webprotege.common.ProjectId;

import javax.annotation.Nonnull;
import java.util.List;

/**
* A {@link PackagedProjectChangeEvent} enriched with the per-project sequence ordinal assigned to
* the bundle when it was durably archived. Published post-persistence on {@link #CHANNEL} so the
* gateway can push it with truthful event tags ({@code startTag = sequenceNumber - 1},
* {@code endTag = sequenceNumber}); because it is emitted only after the archive write, anything
* pushed is already fetchable via the pull path.
* <p>
* The field names are the wire contract consumed by the gateway's mirror DTO. {@code projectId},
* {@code eventId} and {@code projectEvents} serialize exactly as in {@link PackagedProjectChangeEvent}.
*/
@JsonTypeName(SequencedPackagedProjectChangeEvent.CHANNEL)
public record SequencedPackagedProjectChangeEvent(ProjectId projectId,
EventId eventId,
int sequenceNumber,
List<ProjectEvent> projectEvents) implements ProjectEvent {

public static final String CHANNEL = "webprotege.events.projects.SequencedPackagedProjectChange";

@Nonnull
@Override
public ProjectId projectId() {
return projectId;
}

@Nonnull
@Override
public EventId eventId() {
return eventId;
}

@Override
public String getChannel() {
return CHANNEL;
}
}

This file was deleted.

This file was deleted.

Original file line number Diff line number Diff line change
@@ -0,0 +1,39 @@
package edu.stanford.protege.webprotegeeventshistory.sequence;

import org.springframework.data.annotation.Id;
import org.springframework.data.mongodb.core.mapping.Document;

/**
* The per-project sequence counter. One document per project, its {@code _id} being the
* project id and {@code seq} the highest sequence ordinal assigned so far for that project.
* The counter is advanced atomically server-side (single-document {@code $inc}) by
* {@link ProjectSequenceService}.
*/
@Document(collection = ProjectEventSequence.COLLECTION)
public class ProjectEventSequence {

public static final String COLLECTION = "ProjectEventSequence";

public static final String SEQ_FIELD = "seq";

@Id
private String projectId;

private int seq;

public ProjectEventSequence() {
}

public ProjectEventSequence(String projectId, int seq) {
this.projectId = projectId;
this.seq = seq;
}

public String getProjectId() {
return projectId;
}

public int getSeq() {
return seq;
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,74 @@
package edu.stanford.protege.webprotegeeventshistory.sequence;

import edu.stanford.protege.webprotegeeventshistory.HighLevelBusinessEvent;
import org.bson.Document;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.boot.context.event.ApplicationReadyEvent;
import org.springframework.context.event.EventListener;
import org.springframework.data.mongodb.core.MongoTemplate;
import org.springframework.data.mongodb.core.aggregation.Aggregation;
import org.springframework.data.mongodb.core.aggregation.AggregationResults;
import org.springframework.data.mongodb.core.query.Criteria;
import org.springframework.data.mongodb.core.query.Query;
import org.springframework.data.mongodb.core.query.Update;
import org.springframework.stereotype.Component;

/**
* Seeds the per-project {@link ProjectEventSequence} counters from the events already archived
* in {@code HighLevelEvents}. For each project it sets the counter to {@code max(timeStamp)} of
* that project's archived events so newly minted ordinals stay strictly above every bookmark a
* connected client may still hold.
* <p>
* The seed is idempotent: it uses {@code $setOnInsert}, so it only writes a counter that does not
* exist yet and never overwrites one that has already advanced. That makes it safe to run on
* every (possibly rolling) restart. It runs once the application is ready, before the counter is
* consulted on the live traffic path.
*/
@Component
public class ProjectSequenceSeedMigration {

private static final Logger LOGGER = LoggerFactory.getLogger(ProjectSequenceSeedMigration.class);

private static final String GROUP_ID = "_id";

private static final String MAX_SEQ = "maxSeq";

private final MongoTemplate mongoTemplate;

public ProjectSequenceSeedMigration(MongoTemplate mongoTemplate) {
this.mongoTemplate = mongoTemplate;
}

@EventListener(ApplicationReadyEvent.class)
public void onApplicationReady() {
seed();
}

/**
* Inserts a counter set to the per-project archived maximum for every project that does not
* already have one. Existing counters are left untouched.
*/
public void seed() {
var aggregation = Aggregation.newAggregation(
Aggregation.group("projectId").max("timeStamp").as(MAX_SEQ));
AggregationResults<Document> results =
mongoTemplate.aggregate(aggregation, HighLevelBusinessEvent.class, Document.class);

int seeded = 0;
for (Document row : results) {
String projectId = row.getString(GROUP_ID);
if (projectId == null) {
continue;
}
Number maxSeq = (Number) row.get(MAX_SEQ);
int seedValue = maxSeq == null ? 0 : maxSeq.intValue();
mongoTemplate.upsert(
new Query(Criteria.where("_id").is(projectId)),
new Update().setOnInsert(ProjectEventSequence.SEQ_FIELD, seedValue),
ProjectEventSequence.class);
seeded++;
}
LOGGER.info("Project sequence seed migration processed {} project(s)", seeded);
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,49 @@
package edu.stanford.protege.webprotegeeventshistory.sequence;

import org.springframework.data.mongodb.core.FindAndModifyOptions;
import org.springframework.data.mongodb.core.MongoTemplate;
import org.springframework.data.mongodb.core.query.Criteria;
import org.springframework.data.mongodb.core.query.Query;
import org.springframework.data.mongodb.core.query.Update;
import org.springframework.stereotype.Service;

/**
* Assigns per-project, monotonically increasing sequence ordinals.
* <p>
* Each project has its own counter document in the {@code ProjectEventSequence} collection
* (see {@link ProjectEventSequence}). {@link #next(String)} advances it with a single-document
* {@code findAndModify}/{@code $inc}, which MongoDB executes atomically on the server, so the
* numbers stay dense and gapless even when several service instances increment the same
* project's counter concurrently. The first ordinal handed out for a project is {@code 1}.
*/
@Service
public class ProjectSequenceService {

private final MongoTemplate mongoTemplate;

public ProjectSequenceService(MongoTemplate mongoTemplate) {
this.mongoTemplate = mongoTemplate;
}

/**
* Atomically increments and returns the next sequence ordinal for the given project.
* If the project has no counter yet one is created and {@code 1} is returned.
*/
public int next(String projectId) {
var query = new Query(Criteria.where("_id").is(projectId));
var update = new Update().inc(ProjectEventSequence.SEQ_FIELD, 1);
var options = new FindAndModifyOptions().upsert(true).returnNew(true);
ProjectEventSequence result = mongoTemplate.findAndModify(query, update, options, ProjectEventSequence.class);
return result.getSeq();
}

/**
* Reads the current head ordinal for the given project without advancing it.
* Returns {@code 0} when the project has no counter yet (consumed by #301).
*/
public int getCurrentSequence(String projectId) {
var query = new Query(Criteria.where("_id").is(projectId));
ProjectEventSequence result = mongoTemplate.findOne(query, ProjectEventSequence.class);
return result == null ? 0 : result.getSeq();
}
}

This file was deleted.

Loading
Loading