mirror of
https://gitee.com/mateos/mateclaw.git
synced 2026-09-13 03:13:41 +08:00
feat(wiki): page-type count-threshold trigger for pipelines
This commit is contained in:
parent
5278568594
commit
85f5df394b
@ -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.
|
||||
*
|
||||
* <p>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<WikiPipelineDefinitionEntity> defs = definitionMapper.selectList(
|
||||
new LambdaQueryWrapper<WikiPipelineDefinitionEntity>()
|
||||
.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<WikiPageEntity>()
|
||||
.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;
|
||||
}
|
||||
}
|
||||
@ -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.<WikiPipelineRunEntity>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));
|
||||
}
|
||||
}
|
||||
Loading…
Reference in New Issue
Block a user