diff --git a/mateclaw-server/src/main/java/vip/mate/wiki/pipeline/WikiPipelineTriggerService.java b/mateclaw-server/src/main/java/vip/mate/wiki/pipeline/WikiPipelineTriggerService.java new file mode 100644 index 00000000..4f29ba73 --- /dev/null +++ b/mateclaw-server/src/main/java/vip/mate/wiki/pipeline/WikiPipelineTriggerService.java @@ -0,0 +1,114 @@ +package vip.mate.wiki.pipeline; + +import com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper; +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; +import lombok.extern.slf4j.Slf4j; +import org.springframework.stereotype.Service; +import vip.mate.wiki.model.WikiPageEntity; +import vip.mate.wiki.model.WikiPipelineDefinitionEntity; +import vip.mate.wiki.repository.WikiPageMapper; +import vip.mate.wiki.repository.WikiPipelineDefinitionMapper; + +import java.util.List; + +/** + * Evaluates {@code page_type_count} pipeline triggers: when a KB accumulates a + * multiple of the configured threshold of a given pageType, the matching + * pipeline definitions fire once per threshold bucket. + * + *

The bucket is {@code count / threshold}; the run table's unique key makes + * each bucket fire at most once, so re-evaluating on every page create is safe + * and idempotent across instances. + * + * @author MateClaw Team + */ +@Slf4j +@Service +public class WikiPipelineTriggerService { + + private static final String TRIGGER_PAGE_TYPE_COUNT = "page_type_count"; + + private final WikiPipelineDefinitionMapper definitionMapper; + private final WikiPipelineService pipelineService; + private final WikiPageMapper pageMapper; + private final ObjectMapper objectMapper; + + public WikiPipelineTriggerService(WikiPipelineDefinitionMapper definitionMapper, + WikiPipelineService pipelineService, + WikiPageMapper pageMapper, + ObjectMapper objectMapper) { + this.definitionMapper = definitionMapper; + this.pipelineService = pipelineService; + this.pageMapper = pageMapper; + this.objectMapper = objectMapper; + } + + /** + * Re-evaluate count-threshold pipelines for a KB / pageType. Returns the + * number of runs actually started (0 when no threshold bucket was newly + * crossed). Safe to call after every page create. + */ + public int onPageTypeCount(Long kbId, String pageType) { + if (kbId == null || pageType == null || pageType.isBlank()) { + return 0; + } + List defs = definitionMapper.selectList( + new LambdaQueryWrapper() + .eq(WikiPipelineDefinitionEntity::getKbId, kbId) + .eq(WikiPipelineDefinitionEntity::getTriggerType, TRIGGER_PAGE_TYPE_COUNT) + .eq(WikiPipelineDefinitionEntity::getEnabled, 1)); + if (defs.isEmpty()) { + return 0; + } + int started = 0; + for (WikiPipelineDefinitionEntity def : defs) { + TriggerConfig cfg = parseConfig(def.getTriggerConfigJson()); + if (cfg == null || cfg.threshold <= 0 || !pageType.equalsIgnoreCase(cfg.pageType)) { + continue; + } + long count = countPagesOfType(kbId, cfg.pageType); + long bucket = count / cfg.threshold; + if (bucket < 1) { + continue; // threshold not reached yet + } + String input = "{\"pageType\":\"" + cfg.pageType + "\",\"count\":" + count + "}"; + WikiPipelineService.RunOutcome outcome = + pipelineService.execute(def, cfg.pageType, String.valueOf(bucket), input); + if (!outcome.duplicate() && outcome.run() != null) { + started++; + log.info("[WikiPipeline] trigger fired: def={} pageType={} count={} bucket={}", + def.getId(), cfg.pageType, count, bucket); + } + } + return started; + } + + private long countPagesOfType(Long kbId, String pageType) { + return pageMapper.selectCount(new LambdaQueryWrapper() + .eq(WikiPageEntity::getKbId, kbId) + .eq(WikiPageEntity::getPageType, pageType.toLowerCase()) + .and(w -> w.ne(WikiPageEntity::getArchived, 1).or().isNull(WikiPageEntity::getArchived))); + } + + private TriggerConfig parseConfig(String json) { + if (json == null || json.isBlank()) { + return null; + } + try { + JsonNode node = objectMapper.readTree(json); + TriggerConfig cfg = new TriggerConfig(); + cfg.pageType = node.path("page_type").asText(null); + cfg.threshold = node.path("threshold").asInt(0); + return cfg.pageType == null ? null : cfg; + } catch (Exception e) { + log.warn("[WikiPipeline] bad trigger config: {}", e.getMessage()); + return null; + } + } + + private static final class TriggerConfig { + private String pageType; + private int threshold; + } +} diff --git a/mateclaw-server/src/test/java/vip/mate/wiki/pipeline/WikiPipelineTriggerServiceE2ETest.java b/mateclaw-server/src/test/java/vip/mate/wiki/pipeline/WikiPipelineTriggerServiceE2ETest.java new file mode 100644 index 00000000..bb52f312 --- /dev/null +++ b/mateclaw-server/src/test/java/vip/mate/wiki/pipeline/WikiPipelineTriggerServiceE2ETest.java @@ -0,0 +1,123 @@ +package vip.mate.wiki.pipeline; + +import com.baomidou.mybatisplus.core.toolkit.Wrappers; +import com.fasterxml.jackson.databind.ObjectMapper; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.boot.test.context.SpringBootTest; +import vip.mate.wiki.model.WikiPipelineDefinitionEntity; +import vip.mate.wiki.model.WikiPipelineRunEntity; +import vip.mate.wiki.repository.WikiPageMapper; +import vip.mate.wiki.repository.WikiPipelineDefinitionMapper; +import vip.mate.wiki.repository.WikiPipelineRunMapper; +import vip.mate.wiki.repository.WikiPipelineStepRunMapper; +import vip.mate.wiki.service.WikiPageService; + +import java.time.LocalDateTime; +import java.util.List; + +import static org.junit.jupiter.api.Assertions.assertEquals; + +/** + * End-to-end test of the count-threshold trigger against H2: a pipeline fires + * once the page count reaches the threshold, is deduplicated within the same + * threshold bucket, and fires again at the next bucket. + */ +@SpringBootTest( + webEnvironment = SpringBootTest.WebEnvironment.NONE, + properties = { + "spring.flyway.enabled=true", + "spring.flyway.locations=classpath:db/migration/h2", + "mateclaw.feature-flag.refresh-ms=999999" + } +) +class WikiPipelineTriggerServiceE2ETest { + + private static final WikiStepExecutor NOOP = new WikiStepExecutor() { + public String type() { return "noop"; } + public String execute(WikiStepContext c) { return "ok"; } + }; + + @Autowired private WikiPipelineDefinitionMapper definitionMapper; + @Autowired private WikiPipelineRunMapper runMapper; + @Autowired private WikiPipelineStepRunMapper stepRunMapper; + @Autowired private WikiPageMapper pageMapper; + @Autowired private WikiPageService pageService; + @Autowired private ObjectMapper objectMapper; + + private WikiPipelineTriggerService triggerService; + + private static final java.util.concurrent.atomic.AtomicLong SEQ = + new java.util.concurrent.atomic.AtomicLong(System.nanoTime()); + + @BeforeEach + void setUp() { + WikiPipelineService pipelineService = new WikiPipelineService( + runMapper, stepRunMapper, objectMapper, List.of(NOOP)); + triggerService = new WikiPipelineTriggerService( + definitionMapper, pipelineService, pageMapper, objectMapper); + } + + private long seedDefinition(long kb, int threshold) { + WikiPipelineDefinitionEntity d = new WikiPipelineDefinitionEntity(); + d.setKbId(kb); + d.setName("episode-to-pattern-" + SEQ.incrementAndGet()); + d.setOwnerAgentId(42L); + d.setTriggerType("page_type_count"); + d.setTriggerConfigJson("{\"page_type\":\"episode\",\"threshold\":" + threshold + "}"); + d.setStepsJson("[{\"id\":\"s\",\"executor\":\"noop\"}]"); + d.setEnabled(1); + d.setCreateTime(LocalDateTime.now()); + d.setUpdateTime(LocalDateTime.now()); + definitionMapper.insert(d); + return d.getId(); + } + + private void addEpisodes(long kb, int n) { + for (int i = 0; i < n; i++) { + pageService.createPage(kb, "ep-" + kb + "-" + SEQ.incrementAndGet(), + "Episode", "body", "s", "[1]", "episode"); + } + } + + private long runCount(long defId) { + return runMapper.selectCount(Wrappers.lambdaQuery() + .eq(WikiPipelineRunEntity::getDefinitionId, defId)); + } + + @Test + void firesAtThreshold_dedupsWithinBucket_firesAtNextBucket() { + long kb = SEQ.incrementAndGet(); + long defId = seedDefinition(kb, 3); + + // Below threshold: no run. + addEpisodes(kb, 2); + assertEquals(0, triggerService.onPageTypeCount(kb, "episode")); + assertEquals(0, runCount(defId)); + + // Reaching threshold (3) fires one run (bucket 1). + addEpisodes(kb, 1); + assertEquals(1, triggerService.onPageTypeCount(kb, "episode")); + assertEquals(1, runCount(defId)); + + // Still in bucket 1 (count 4): deduped, no new run. + addEpisodes(kb, 1); + assertEquals(0, triggerService.onPageTypeCount(kb, "episode")); + assertEquals(1, runCount(defId)); + + // Crossing into bucket 2 (count 6) fires again. + addEpisodes(kb, 2); + assertEquals(1, triggerService.onPageTypeCount(kb, "episode")); + assertEquals(2, runCount(defId)); + } + + @Test + void nonMatchingPageType_doesNotFire() { + long kb = SEQ.incrementAndGet(); + long defId = seedDefinition(kb, 1); + // A different pageType event must not trigger the episode pipeline. + assertEquals(0, triggerService.onPageTypeCount(kb, "concept")); + assertEquals(0, runCount(defId)); + } +}