feat(subscription): notify subscribers on skill publish and version yank

This commit is contained in:
dongmucat 2026-04-27 16:51:40 +08:00
parent 6f0103f013
commit 9bc5d0c784

View file

@ -7,6 +7,7 @@ import com.iflytek.skillhub.domain.namespace.NamespaceRepository;
import com.iflytek.skillhub.domain.skill.Skill;
import com.iflytek.skillhub.domain.skill.SkillRepository;
import com.iflytek.skillhub.domain.skill.SkillVersionRepository;
import com.iflytek.skillhub.domain.social.SkillSubscriptionService;
import com.iflytek.skillhub.notification.domain.NotificationCategory;
import com.iflytek.skillhub.notification.service.NotificationDispatcher;
import org.slf4j.Logger;
@ -29,6 +30,7 @@ public class NotificationEventListener {
private final NamespaceRepository namespaceRepository;
private final RecipientResolver recipientResolver;
private final NotificationDispatcher dispatcher;
private final SkillSubscriptionService skillSubscriptionService;
private final ObjectMapper objectMapper;
public NotificationEventListener(SkillRepository skillRepository,
@ -36,12 +38,14 @@ public class NotificationEventListener {
NamespaceRepository namespaceRepository,
RecipientResolver recipientResolver,
NotificationDispatcher dispatcher,
SkillSubscriptionService skillSubscriptionService,
ObjectMapper objectMapper) {
this.skillRepository = skillRepository;
this.skillVersionRepository = skillVersionRepository;
this.namespaceRepository = namespaceRepository;
this.recipientResolver = recipientResolver;
this.dispatcher = dispatcher;
this.skillSubscriptionService = skillSubscriptionService;
this.objectMapper = objectMapper;
}
@ -61,6 +65,50 @@ public class NotificationEventListener {
});
}
@Async("skillhubEventExecutor")
@TransactionalEventListener
public void onSkillPublishedForSubscribers(SkillPublishedEvent event) {
skillRepository.findById(event.skillId()).ifPresent(skill -> {
List<String> subscribers = skillSubscriptionService.findSubscribersBySkillId(event.skillId());
if (subscribers.isEmpty()) {
return;
}
String title = "Skill updated: " + skillDisplayName(skill);
Map<String, Object> body = bodyWithSkill(skill);
versionLabel(event.versionId(), body);
String json = toJson(body);
for (String subscriberId : subscribers) {
if (subscriberId.equals(event.publisherId())) {
continue; // skip the publisher
}
dispatcher.dispatch(subscriberId, NotificationCategory.PUBLISH,
"SUBSCRIPTION_NEW_VERSION", title, json, "SKILL", event.skillId());
}
});
}
@Async("skillhubEventExecutor")
@TransactionalEventListener
public void onSkillVersionYankedForSubscribers(SkillVersionYankedEvent event) {
skillRepository.findById(event.skillId()).ifPresent(skill -> {
List<String> subscribers = skillSubscriptionService.findSubscribersBySkillId(event.skillId());
if (subscribers.isEmpty()) {
return;
}
String title = "Skill version yanked: " + skillDisplayName(skill);
Map<String, Object> body = bodyWithSkill(skill);
versionLabel(event.versionId(), body);
String json = toJson(body);
for (String subscriberId : subscribers) {
if (subscriberId.equals(event.actorUserId())) {
continue; // skip the actor
}
dispatcher.dispatch(subscriberId, NotificationCategory.PUBLISH,
"SUBSCRIPTION_VERSION_YANKED", title, json, "SKILL", event.skillId());
}
});
}
@Async("skillhubEventExecutor")
@TransactionalEventListener
public void onReviewSubmitted(ReviewSubmittedEvent event) {