调整项目结构
This commit is contained in:
@@ -0,0 +1,239 @@
|
||||
package tech.easyflow.manuagent;
|
||||
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
import static org.mockito.Mockito.mock;
|
||||
import static org.mockito.Mockito.when;
|
||||
|
||||
import tech.easyflow.manuagent.agent.AgentEventService;
|
||||
import tech.easyflow.manuagent.artifact.ArtifactService;
|
||||
import tech.easyflow.manuagent.artifact.DocxValidator;
|
||||
import tech.easyflow.manuagent.auth.UserService;
|
||||
import tech.easyflow.manuagent.project.ProjectService;
|
||||
import tech.easyflow.manuagent.project.ProjectFileService;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.sql.DriverManager;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.UUID;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import org.flywaydb.core.Flyway;
|
||||
import org.junit.jupiter.api.BeforeAll;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
import org.postgresql.ds.PGSimpleDataSource;
|
||||
import org.springframework.jdbc.core.simple.JdbcClient;
|
||||
import org.testcontainers.containers.PostgreSQLContainer;
|
||||
import org.testcontainers.junit.jupiter.Container;
|
||||
import org.testcontainers.junit.jupiter.Testcontainers;
|
||||
|
||||
/**
|
||||
* 验证 PostgreSQL 17 全量迁移及事件游标回放。
|
||||
*/
|
||||
@Testcontainers
|
||||
class DatabaseAndEventIntegrationTest {
|
||||
|
||||
@Container
|
||||
private static final PostgreSQLContainer<?> POSTGRES = new PostgreSQLContainer<>("postgres:17-alpine");
|
||||
|
||||
@TempDir
|
||||
private Path temporaryDirectory;
|
||||
|
||||
/**
|
||||
* 在干净 PostgreSQL 17 实例执行并校验全部 Flyway 迁移。
|
||||
*/
|
||||
@BeforeAll
|
||||
static void migrate() {
|
||||
Flyway flyway = Flyway.configure()
|
||||
.dataSource(POSTGRES.getJdbcUrl(), POSTGRES.getUsername(), POSTGRES.getPassword())
|
||||
.locations("classpath:db/migration")
|
||||
.load();
|
||||
flyway.migrate();
|
||||
flyway.validate();
|
||||
assertThat(flyway.migrate().migrationsExecuted).isZero();
|
||||
}
|
||||
|
||||
/**
|
||||
* 验证核心表、索引和约束已建立。
|
||||
*
|
||||
* @throws Exception 数据库访问失败时抛出
|
||||
*/
|
||||
@Test
|
||||
void shouldCreateCoreSchemaOnPostgres17() throws Exception {
|
||||
try (var connection = DriverManager.getConnection(
|
||||
POSTGRES.getJdbcUrl(), POSTGRES.getUsername(), POSTGRES.getPassword());
|
||||
var statement = connection.createStatement();
|
||||
var result = statement.executeQuery("""
|
||||
SELECT count(*) FROM information_schema.tables
|
||||
WHERE table_schema IN ('app', 'agentscope')
|
||||
""")) {
|
||||
assertThat(result.next()).isTrue();
|
||||
assertThat(result.getInt(1)).isEqualTo(12);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 验证事件按项目全局 ID 增量回放且不重复。
|
||||
*/
|
||||
@Test
|
||||
void shouldReplayEventsAfterCursorInOrder() {
|
||||
JdbcClient jdbc = jdbc();
|
||||
UUID userId = UUID.randomUUID();
|
||||
UUID modelId = UUID.randomUUID();
|
||||
UUID projectId = UUID.randomUUID();
|
||||
UUID runId = UUID.randomUUID();
|
||||
jdbc.sql("INSERT INTO app.app_user(id, username, password_hash, display_name) VALUES (:id, :name, 'x', 'test')")
|
||||
.param("id", userId).param("name", "u-" + userId).update();
|
||||
jdbc.sql("""
|
||||
INSERT INTO app.model_config(id, name, provider, base_url, model_id, is_default)
|
||||
VALUES (:id, :name, 'OPENAI_COMPATIBLE', 'https://example.test', 'model', TRUE)
|
||||
""").param("id", modelId).param("name", "m-" + modelId).update();
|
||||
jdbc.sql("""
|
||||
INSERT INTO app.project(id, company_name, project_name, agui_thread_id, application_level, created_by)
|
||||
VALUES (:id, '企业', '项目', :thread, 'ADVANCED', :userId)
|
||||
""").param("id", projectId).param("thread", "t-" + projectId).param("userId", userId).update();
|
||||
jdbc.sql("""
|
||||
INSERT INTO app.agent_run(id, project_id, model_config_id, trigger_type, status, trace_id)
|
||||
VALUES (:id, :projectId, :modelId, 'INITIAL', 'RUNNING', :trace)
|
||||
""").param("id", runId).param("projectId", projectId).param("modelId", modelId)
|
||||
.param("trace", UUID.randomUUID().toString()).update();
|
||||
|
||||
AgentEventService service = new AgentEventService(jdbc, new ObjectMapper());
|
||||
long first = service.append(projectId, runId, "RUN_STARTED", Map.of("phase", "MATERIAL_CHECK")).id();
|
||||
long second = service.append(projectId, runId, "TEXT_MESSAGE_CONTENT", Map.of("delta", "分析")).id();
|
||||
long third = service.append(projectId, runId, "TEXT_MESSAGE_CONTENT", Map.of("delta", "完成")).id();
|
||||
|
||||
assertThat(service.listAfter(projectId, first, 100))
|
||||
.extracting(AgentEventService.EventView::id)
|
||||
.containsExactly(second, third);
|
||||
|
||||
var next = service.streamAfter(projectId, third)
|
||||
.filter(event -> !"HEARTBEAT".equals(event.type()))
|
||||
.next()
|
||||
.toFuture();
|
||||
long pushed = service.append(projectId, runId, "TEXT_MESSAGE_CONTENT", Map.of("delta", "推送")).id();
|
||||
assertThat(next.orTimeout(2, TimeUnit.SECONDS).join().id()).isEqualTo(pushed);
|
||||
}
|
||||
|
||||
/**
|
||||
* 验证重试生成同一路径产物时更新登记信息,避免唯一约束导致成功 Run 被标记失败。
|
||||
*
|
||||
* @throws Exception 临时文件写入失败时抛出
|
||||
*/
|
||||
@Test
|
||||
void shouldReplaceArtifactMetadataForSameProjectPath() throws Exception {
|
||||
JdbcClient jdbc = jdbc();
|
||||
UUID userId = UUID.randomUUID();
|
||||
UUID modelId = UUID.randomUUID();
|
||||
UUID projectId = UUID.randomUUID();
|
||||
UUID firstRunId = UUID.randomUUID();
|
||||
UUID secondRunId = UUID.randomUUID();
|
||||
jdbc.sql("INSERT INTO app.app_user(id, username, password_hash, display_name) VALUES (:id, :name, 'x', 'test')")
|
||||
.param("id", userId).param("name", "u-" + userId).update();
|
||||
jdbc.sql("""
|
||||
INSERT INTO app.model_config(id, name, provider, base_url, model_id)
|
||||
VALUES (:id, :name, 'OPENAI_COMPATIBLE', 'https://example.test', 'model')
|
||||
""").param("id", modelId).param("name", "m-" + modelId).update();
|
||||
jdbc.sql("""
|
||||
INSERT INTO app.project(id, company_name, project_name, agui_thread_id, application_level, created_by)
|
||||
VALUES (:id, '企业', '项目', :thread, 'ADVANCED', :userId)
|
||||
""").param("id", projectId).param("thread", "t-" + projectId).param("userId", userId).update();
|
||||
for (UUID runId : List.of(firstRunId, secondRunId)) {
|
||||
jdbc.sql("""
|
||||
INSERT INTO app.agent_run(
|
||||
id, project_id, model_config_id, trigger_type, status, trace_id, ended_at)
|
||||
VALUES (:id, :projectId, :modelId, 'RETRY', 'COMPLETED', :trace, CURRENT_TIMESTAMP)
|
||||
""").param("id", runId).param("projectId", projectId).param("modelId", modelId)
|
||||
.param("trace", UUID.randomUUID().toString()).update();
|
||||
}
|
||||
|
||||
Path document = temporaryDirectory.resolve("draft.docx");
|
||||
Files.writeString(document, "first");
|
||||
ProjectFileService files = mock(ProjectFileService.class);
|
||||
when(files.safeProjectPath(projectId, "artifacts/draft.docx")).thenReturn(document);
|
||||
ArtifactService artifacts = new ArtifactService(jdbc, files, new DocxValidator());
|
||||
ObjectMapper mapper = new ObjectMapper();
|
||||
|
||||
ArtifactService.ArtifactView first = artifacts.publish(
|
||||
projectId, firstRunId, "DOCX", "draft.docx", "artifacts/draft.docx",
|
||||
mapper.createObjectNode().put("version", 1));
|
||||
Files.writeString(document, "second version");
|
||||
ArtifactService.ArtifactView second = artifacts.publish(
|
||||
projectId, secondRunId, "DOCX", "draft.docx", "artifacts/draft.docx",
|
||||
mapper.createObjectNode().put("version", 2));
|
||||
|
||||
assertThat(second.id()).isEqualTo(first.id());
|
||||
assertThat(second.runId()).isEqualTo(secondRunId);
|
||||
assertThat(second.sizeBytes()).isEqualTo(Files.size(document));
|
||||
assertThat(artifacts.list(projectId)).hasSize(1);
|
||||
}
|
||||
|
||||
/**
|
||||
* 验证项目真删除会清除所有关联业务记录。
|
||||
*/
|
||||
@Test
|
||||
void shouldDeleteProjectRecords() {
|
||||
JdbcClient jdbc = jdbc();
|
||||
UUID userId = UUID.randomUUID();
|
||||
UUID projectId = UUID.randomUUID();
|
||||
UUID runId = UUID.randomUUID();
|
||||
jdbc.sql("INSERT INTO app.app_user(id, username, password_hash, display_name) VALUES (:id, :name, 'x', 'test')")
|
||||
.param("id", userId).param("name", "u-" + userId).update();
|
||||
jdbc.sql("""
|
||||
INSERT INTO app.project(id, company_name, project_name, agui_thread_id, application_level, created_by)
|
||||
VALUES (:id, '待删除企业', '待删除项目', :thread, 'ADVANCED', :userId)
|
||||
""").param("id", projectId).param("thread", "t-" + projectId).param("userId", userId).update();
|
||||
jdbc.sql("""
|
||||
INSERT INTO app.agent_run(id, project_id, trigger_type, status, trace_id, ended_at)
|
||||
VALUES (:id, :projectId, 'INITIAL', 'COMPLETED', :trace, CURRENT_TIMESTAMP)
|
||||
""").param("id", runId).param("projectId", projectId)
|
||||
.param("trace", UUID.randomUUID().toString()).update();
|
||||
jdbc.sql("""
|
||||
INSERT INTO app.agent_event(project_id, run_id, event_type, payload)
|
||||
VALUES (:projectId, :runId, 'RUN_FINISHED', '{}'::jsonb)
|
||||
""").param("projectId", projectId).param("runId", runId).update();
|
||||
jdbc.sql("""
|
||||
INSERT INTO app.project_plan(id, project_id, plan_version, status, plan_json, created_by)
|
||||
VALUES (:id, :projectId, 1, 'DRAFT', '{}'::jsonb, :userId)
|
||||
""").param("id", UUID.randomUUID()).param("projectId", projectId).param("userId", userId).update();
|
||||
jdbc.sql("""
|
||||
INSERT INTO app.project_file(
|
||||
id, project_id, original_name, stored_name, relative_path, mime_type,
|
||||
extension, size_bytes, sha256, uploaded_by)
|
||||
VALUES (:id, :projectId, 'input.txt', 'input.txt', 'inputs/input.txt',
|
||||
'text/plain', 'txt', 1, :sha, :userId)
|
||||
""").param("id", UUID.randomUUID()).param("projectId", projectId)
|
||||
.param("sha", "0".repeat(64)).param("userId", userId).update();
|
||||
jdbc.sql("""
|
||||
INSERT INTO app.artifact(
|
||||
id, project_id, run_id, kind, name, relative_path, mime_type, size_bytes, sha256)
|
||||
VALUES (:id, :projectId, :runId, 'OTHER', 'result.txt', 'artifacts/result.txt',
|
||||
'text/plain', 1, :sha)
|
||||
""").param("id", UUID.randomUUID()).param("projectId", projectId).param("runId", runId)
|
||||
.param("sha", "0".repeat(64)).update();
|
||||
|
||||
ProjectService service = new ProjectService(jdbc, mock(UserService.class), new ObjectMapper());
|
||||
service.delete(projectId);
|
||||
|
||||
for (String table : List.of("agent_event", "artifact", "project_plan", "project_file", "agent_run")) {
|
||||
Long count = jdbc.sql("SELECT COUNT(*) FROM app." + table + " WHERE project_id = :projectId")
|
||||
.param("projectId", projectId)
|
||||
.query(Long.class)
|
||||
.single();
|
||||
assertThat(count).as(table).isZero();
|
||||
}
|
||||
assertThat(jdbc.sql("SELECT COUNT(*) FROM app.project WHERE id = :projectId")
|
||||
.param("projectId", projectId)
|
||||
.query(Long.class)
|
||||
.single()).isZero();
|
||||
}
|
||||
|
||||
private JdbcClient jdbc() {
|
||||
PGSimpleDataSource source = new PGSimpleDataSource();
|
||||
source.setURL(POSTGRES.getJdbcUrl());
|
||||
source.setUser(POSTGRES.getUsername());
|
||||
source.setPassword(POSTGRES.getPassword());
|
||||
return JdbcClient.create(source);
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user