Compare commits
1 Commits
feat/ms23-
...
feat/ms23-
| Author | SHA1 | Date | |
|---|---|---|---|
| 47ea93ca82 |
@@ -25,14 +25,14 @@ export class AgentIngestionService {
|
|||||||
where: { sessionId: agentId },
|
where: { sessionId: agentId },
|
||||||
create: {
|
create: {
|
||||||
sessionId: agentId,
|
sessionId: agentId,
|
||||||
parentSessionId: parentAgentId ?? null,
|
parentSessionId: parentAgentId,
|
||||||
missionId,
|
missionId,
|
||||||
taskId,
|
taskId,
|
||||||
agentType,
|
agentType,
|
||||||
status: "spawning",
|
status: "spawning",
|
||||||
},
|
},
|
||||||
update: {
|
update: {
|
||||||
parentSessionId: parentAgentId ?? null,
|
parentSessionId: parentAgentId,
|
||||||
missionId,
|
missionId,
|
||||||
taskId,
|
taskId,
|
||||||
agentType,
|
agentType,
|
||||||
|
|||||||
@@ -1,54 +0,0 @@
|
|||||||
import {
|
|
||||||
Body,
|
|
||||||
Controller,
|
|
||||||
Delete,
|
|
||||||
Get,
|
|
||||||
Param,
|
|
||||||
Patch,
|
|
||||||
Post,
|
|
||||||
UseGuards,
|
|
||||||
UsePipes,
|
|
||||||
ValidationPipe,
|
|
||||||
} from "@nestjs/common";
|
|
||||||
import type { AgentProviderConfig } from "@prisma/client";
|
|
||||||
import { OrchestratorApiKeyGuard } from "../../common/guards/api-key.guard";
|
|
||||||
import { OrchestratorThrottlerGuard } from "../../common/guards/throttler.guard";
|
|
||||||
import { AgentProvidersService } from "./agent-providers.service";
|
|
||||||
import { CreateAgentProviderDto } from "./dto/create-agent-provider.dto";
|
|
||||||
import { UpdateAgentProviderDto } from "./dto/update-agent-provider.dto";
|
|
||||||
|
|
||||||
@Controller("agent-providers")
|
|
||||||
@UseGuards(OrchestratorApiKeyGuard, OrchestratorThrottlerGuard)
|
|
||||||
export class AgentProvidersController {
|
|
||||||
constructor(private readonly agentProvidersService: AgentProvidersService) {}
|
|
||||||
|
|
||||||
@Get()
|
|
||||||
async list(): Promise<AgentProviderConfig[]> {
|
|
||||||
return this.agentProvidersService.list();
|
|
||||||
}
|
|
||||||
|
|
||||||
@Get(":id")
|
|
||||||
async getById(@Param("id") id: string): Promise<AgentProviderConfig> {
|
|
||||||
return this.agentProvidersService.getById(id);
|
|
||||||
}
|
|
||||||
|
|
||||||
@Post()
|
|
||||||
@UsePipes(new ValidationPipe({ transform: true, whitelist: true }))
|
|
||||||
async create(@Body() dto: CreateAgentProviderDto): Promise<AgentProviderConfig> {
|
|
||||||
return this.agentProvidersService.create(dto);
|
|
||||||
}
|
|
||||||
|
|
||||||
@Patch(":id")
|
|
||||||
@UsePipes(new ValidationPipe({ transform: true, whitelist: true }))
|
|
||||||
async update(
|
|
||||||
@Param("id") id: string,
|
|
||||||
@Body() dto: UpdateAgentProviderDto
|
|
||||||
): Promise<AgentProviderConfig> {
|
|
||||||
return this.agentProvidersService.update(id, dto);
|
|
||||||
}
|
|
||||||
|
|
||||||
@Delete(":id")
|
|
||||||
async delete(@Param("id") id: string): Promise<AgentProviderConfig> {
|
|
||||||
return this.agentProvidersService.delete(id);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -1,12 +0,0 @@
|
|||||||
import { Module } from "@nestjs/common";
|
|
||||||
import { PrismaModule } from "../../prisma/prisma.module";
|
|
||||||
import { OrchestratorApiKeyGuard } from "../../common/guards/api-key.guard";
|
|
||||||
import { AgentProvidersController } from "./agent-providers.controller";
|
|
||||||
import { AgentProvidersService } from "./agent-providers.service";
|
|
||||||
|
|
||||||
@Module({
|
|
||||||
imports: [PrismaModule],
|
|
||||||
controllers: [AgentProvidersController],
|
|
||||||
providers: [OrchestratorApiKeyGuard, AgentProvidersService],
|
|
||||||
})
|
|
||||||
export class AgentProvidersModule {}
|
|
||||||
@@ -1,211 +0,0 @@
|
|||||||
import { beforeEach, describe, expect, it, vi } from "vitest";
|
|
||||||
import { NotFoundException } from "@nestjs/common";
|
|
||||||
import { AgentProvidersService } from "./agent-providers.service";
|
|
||||||
import { PrismaService } from "../../prisma/prisma.service";
|
|
||||||
|
|
||||||
describe("AgentProvidersService", () => {
|
|
||||||
let service: AgentProvidersService;
|
|
||||||
let prisma: {
|
|
||||||
agentProviderConfig: {
|
|
||||||
findMany: ReturnType<typeof vi.fn>;
|
|
||||||
findUnique: ReturnType<typeof vi.fn>;
|
|
||||||
create: ReturnType<typeof vi.fn>;
|
|
||||||
update: ReturnType<typeof vi.fn>;
|
|
||||||
delete: ReturnType<typeof vi.fn>;
|
|
||||||
};
|
|
||||||
};
|
|
||||||
|
|
||||||
beforeEach(() => {
|
|
||||||
prisma = {
|
|
||||||
agentProviderConfig: {
|
|
||||||
findMany: vi.fn(),
|
|
||||||
findUnique: vi.fn(),
|
|
||||||
create: vi.fn(),
|
|
||||||
update: vi.fn(),
|
|
||||||
delete: vi.fn(),
|
|
||||||
},
|
|
||||||
};
|
|
||||||
|
|
||||||
service = new AgentProvidersService(prisma as unknown as PrismaService);
|
|
||||||
});
|
|
||||||
|
|
||||||
it("lists all provider configs", async () => {
|
|
||||||
const expected = [
|
|
||||||
{
|
|
||||||
id: "cfg-1",
|
|
||||||
workspaceId: "8bcd7eda-a122-4d6c-adfd-b152f6f75369",
|
|
||||||
name: "Primary",
|
|
||||||
provider: "openai",
|
|
||||||
gatewayUrl: "https://gateway.example.com",
|
|
||||||
credentials: {},
|
|
||||||
isActive: true,
|
|
||||||
createdAt: new Date("2026-03-07T18:00:00.000Z"),
|
|
||||||
updatedAt: new Date("2026-03-07T18:00:00.000Z"),
|
|
||||||
},
|
|
||||||
];
|
|
||||||
prisma.agentProviderConfig.findMany.mockResolvedValue(expected);
|
|
||||||
|
|
||||||
const result = await service.list();
|
|
||||||
|
|
||||||
expect(prisma.agentProviderConfig.findMany).toHaveBeenCalledWith({
|
|
||||||
orderBy: [{ createdAt: "desc" }, { id: "desc" }],
|
|
||||||
});
|
|
||||||
expect(result).toEqual(expected);
|
|
||||||
});
|
|
||||||
|
|
||||||
it("returns a single provider config", async () => {
|
|
||||||
const expected = {
|
|
||||||
id: "cfg-1",
|
|
||||||
workspaceId: "8bcd7eda-a122-4d6c-adfd-b152f6f75369",
|
|
||||||
name: "Primary",
|
|
||||||
provider: "openai",
|
|
||||||
gatewayUrl: "https://gateway.example.com",
|
|
||||||
credentials: { apiKeyRef: "vault:openai" },
|
|
||||||
isActive: true,
|
|
||||||
createdAt: new Date("2026-03-07T18:00:00.000Z"),
|
|
||||||
updatedAt: new Date("2026-03-07T18:00:00.000Z"),
|
|
||||||
};
|
|
||||||
prisma.agentProviderConfig.findUnique.mockResolvedValue(expected);
|
|
||||||
|
|
||||||
const result = await service.getById("cfg-1");
|
|
||||||
|
|
||||||
expect(prisma.agentProviderConfig.findUnique).toHaveBeenCalledWith({
|
|
||||||
where: { id: "cfg-1" },
|
|
||||||
});
|
|
||||||
expect(result).toEqual(expected);
|
|
||||||
});
|
|
||||||
|
|
||||||
it("throws NotFoundException when provider config is missing", async () => {
|
|
||||||
prisma.agentProviderConfig.findUnique.mockResolvedValue(null);
|
|
||||||
|
|
||||||
await expect(service.getById("missing")).rejects.toBeInstanceOf(NotFoundException);
|
|
||||||
});
|
|
||||||
|
|
||||||
it("creates a provider config with default credentials", async () => {
|
|
||||||
const created = {
|
|
||||||
id: "cfg-created",
|
|
||||||
workspaceId: "8bcd7eda-a122-4d6c-adfd-b152f6f75369",
|
|
||||||
name: "New Provider",
|
|
||||||
provider: "claude",
|
|
||||||
gatewayUrl: "https://gateway.example.com",
|
|
||||||
credentials: {},
|
|
||||||
isActive: true,
|
|
||||||
createdAt: new Date("2026-03-07T18:00:00.000Z"),
|
|
||||||
updatedAt: new Date("2026-03-07T18:00:00.000Z"),
|
|
||||||
};
|
|
||||||
prisma.agentProviderConfig.create.mockResolvedValue(created);
|
|
||||||
|
|
||||||
const result = await service.create({
|
|
||||||
workspaceId: "8bcd7eda-a122-4d6c-adfd-b152f6f75369",
|
|
||||||
name: "New Provider",
|
|
||||||
provider: "claude",
|
|
||||||
gatewayUrl: "https://gateway.example.com",
|
|
||||||
});
|
|
||||||
|
|
||||||
expect(prisma.agentProviderConfig.create).toHaveBeenCalledWith({
|
|
||||||
data: {
|
|
||||||
workspaceId: "8bcd7eda-a122-4d6c-adfd-b152f6f75369",
|
|
||||||
name: "New Provider",
|
|
||||||
provider: "claude",
|
|
||||||
gatewayUrl: "https://gateway.example.com",
|
|
||||||
credentials: {},
|
|
||||||
},
|
|
||||||
});
|
|
||||||
expect(result).toEqual(created);
|
|
||||||
});
|
|
||||||
|
|
||||||
it("updates a provider config", async () => {
|
|
||||||
prisma.agentProviderConfig.findUnique.mockResolvedValue({
|
|
||||||
id: "cfg-1",
|
|
||||||
workspaceId: "8bcd7eda-a122-4d6c-adfd-b152f6f75369",
|
|
||||||
name: "Primary",
|
|
||||||
provider: "openai",
|
|
||||||
gatewayUrl: "https://gateway.example.com",
|
|
||||||
credentials: {},
|
|
||||||
isActive: true,
|
|
||||||
createdAt: new Date("2026-03-07T18:00:00.000Z"),
|
|
||||||
updatedAt: new Date("2026-03-07T18:00:00.000Z"),
|
|
||||||
});
|
|
||||||
|
|
||||||
const updated = {
|
|
||||||
id: "cfg-1",
|
|
||||||
workspaceId: "8bcd7eda-a122-4d6c-adfd-b152f6f75369",
|
|
||||||
name: "Secondary",
|
|
||||||
provider: "openai",
|
|
||||||
gatewayUrl: "https://gateway2.example.com",
|
|
||||||
credentials: { apiKeyRef: "vault:new" },
|
|
||||||
isActive: false,
|
|
||||||
createdAt: new Date("2026-03-07T18:00:00.000Z"),
|
|
||||||
updatedAt: new Date("2026-03-07T19:00:00.000Z"),
|
|
||||||
};
|
|
||||||
prisma.agentProviderConfig.update.mockResolvedValue(updated);
|
|
||||||
|
|
||||||
const result = await service.update("cfg-1", {
|
|
||||||
name: "Secondary",
|
|
||||||
gatewayUrl: "https://gateway2.example.com",
|
|
||||||
credentials: { apiKeyRef: "vault:new" },
|
|
||||||
isActive: false,
|
|
||||||
});
|
|
||||||
|
|
||||||
expect(prisma.agentProviderConfig.update).toHaveBeenCalledWith({
|
|
||||||
where: { id: "cfg-1" },
|
|
||||||
data: {
|
|
||||||
name: "Secondary",
|
|
||||||
gatewayUrl: "https://gateway2.example.com",
|
|
||||||
credentials: { apiKeyRef: "vault:new" },
|
|
||||||
isActive: false,
|
|
||||||
},
|
|
||||||
});
|
|
||||||
expect(result).toEqual(updated);
|
|
||||||
});
|
|
||||||
|
|
||||||
it("throws NotFoundException when updating a missing provider config", async () => {
|
|
||||||
prisma.agentProviderConfig.findUnique.mockResolvedValue(null);
|
|
||||||
|
|
||||||
await expect(service.update("missing", { name: "Updated" })).rejects.toBeInstanceOf(
|
|
||||||
NotFoundException
|
|
||||||
);
|
|
||||||
expect(prisma.agentProviderConfig.update).not.toHaveBeenCalled();
|
|
||||||
});
|
|
||||||
|
|
||||||
it("deletes a provider config", async () => {
|
|
||||||
prisma.agentProviderConfig.findUnique.mockResolvedValue({
|
|
||||||
id: "cfg-1",
|
|
||||||
workspaceId: "8bcd7eda-a122-4d6c-adfd-b152f6f75369",
|
|
||||||
name: "Primary",
|
|
||||||
provider: "openai",
|
|
||||||
gatewayUrl: "https://gateway.example.com",
|
|
||||||
credentials: {},
|
|
||||||
isActive: true,
|
|
||||||
createdAt: new Date("2026-03-07T18:00:00.000Z"),
|
|
||||||
updatedAt: new Date("2026-03-07T18:00:00.000Z"),
|
|
||||||
});
|
|
||||||
|
|
||||||
const deleted = {
|
|
||||||
id: "cfg-1",
|
|
||||||
workspaceId: "8bcd7eda-a122-4d6c-adfd-b152f6f75369",
|
|
||||||
name: "Primary",
|
|
||||||
provider: "openai",
|
|
||||||
gatewayUrl: "https://gateway.example.com",
|
|
||||||
credentials: {},
|
|
||||||
isActive: true,
|
|
||||||
createdAt: new Date("2026-03-07T18:00:00.000Z"),
|
|
||||||
updatedAt: new Date("2026-03-07T18:00:00.000Z"),
|
|
||||||
};
|
|
||||||
prisma.agentProviderConfig.delete.mockResolvedValue(deleted);
|
|
||||||
|
|
||||||
const result = await service.delete("cfg-1");
|
|
||||||
|
|
||||||
expect(prisma.agentProviderConfig.delete).toHaveBeenCalledWith({
|
|
||||||
where: { id: "cfg-1" },
|
|
||||||
});
|
|
||||||
expect(result).toEqual(deleted);
|
|
||||||
});
|
|
||||||
|
|
||||||
it("throws NotFoundException when deleting a missing provider config", async () => {
|
|
||||||
prisma.agentProviderConfig.findUnique.mockResolvedValue(null);
|
|
||||||
|
|
||||||
await expect(service.delete("missing")).rejects.toBeInstanceOf(NotFoundException);
|
|
||||||
expect(prisma.agentProviderConfig.delete).not.toHaveBeenCalled();
|
|
||||||
});
|
|
||||||
});
|
|
||||||
@@ -1,71 +0,0 @@
|
|||||||
import { Injectable, NotFoundException } from "@nestjs/common";
|
|
||||||
import type { AgentProviderConfig, Prisma } from "@prisma/client";
|
|
||||||
import { PrismaService } from "../../prisma/prisma.service";
|
|
||||||
import { CreateAgentProviderDto } from "./dto/create-agent-provider.dto";
|
|
||||||
import { UpdateAgentProviderDto } from "./dto/update-agent-provider.dto";
|
|
||||||
|
|
||||||
@Injectable()
|
|
||||||
export class AgentProvidersService {
|
|
||||||
constructor(private readonly prisma: PrismaService) {}
|
|
||||||
|
|
||||||
async list(): Promise<AgentProviderConfig[]> {
|
|
||||||
return this.prisma.agentProviderConfig.findMany({
|
|
||||||
orderBy: [{ createdAt: "desc" }, { id: "desc" }],
|
|
||||||
});
|
|
||||||
}
|
|
||||||
|
|
||||||
async getById(id: string): Promise<AgentProviderConfig> {
|
|
||||||
const providerConfig = await this.prisma.agentProviderConfig.findUnique({
|
|
||||||
where: { id },
|
|
||||||
});
|
|
||||||
|
|
||||||
if (!providerConfig) {
|
|
||||||
throw new NotFoundException(`Agent provider config with id ${id} not found`);
|
|
||||||
}
|
|
||||||
|
|
||||||
return providerConfig;
|
|
||||||
}
|
|
||||||
|
|
||||||
async create(dto: CreateAgentProviderDto): Promise<AgentProviderConfig> {
|
|
||||||
return this.prisma.agentProviderConfig.create({
|
|
||||||
data: {
|
|
||||||
workspaceId: dto.workspaceId,
|
|
||||||
name: dto.name,
|
|
||||||
provider: dto.provider,
|
|
||||||
gatewayUrl: dto.gatewayUrl,
|
|
||||||
credentials: this.toJsonValue(dto.credentials ?? {}),
|
|
||||||
...(dto.isActive !== undefined ? { isActive: dto.isActive } : {}),
|
|
||||||
},
|
|
||||||
});
|
|
||||||
}
|
|
||||||
|
|
||||||
async update(id: string, dto: UpdateAgentProviderDto): Promise<AgentProviderConfig> {
|
|
||||||
await this.getById(id);
|
|
||||||
|
|
||||||
const data: Prisma.AgentProviderConfigUpdateInput = {
|
|
||||||
...(dto.workspaceId !== undefined ? { workspaceId: dto.workspaceId } : {}),
|
|
||||||
...(dto.name !== undefined ? { name: dto.name } : {}),
|
|
||||||
...(dto.provider !== undefined ? { provider: dto.provider } : {}),
|
|
||||||
...(dto.gatewayUrl !== undefined ? { gatewayUrl: dto.gatewayUrl } : {}),
|
|
||||||
...(dto.isActive !== undefined ? { isActive: dto.isActive } : {}),
|
|
||||||
...(dto.credentials !== undefined ? { credentials: this.toJsonValue(dto.credentials) } : {}),
|
|
||||||
};
|
|
||||||
|
|
||||||
return this.prisma.agentProviderConfig.update({
|
|
||||||
where: { id },
|
|
||||||
data,
|
|
||||||
});
|
|
||||||
}
|
|
||||||
|
|
||||||
async delete(id: string): Promise<AgentProviderConfig> {
|
|
||||||
await this.getById(id);
|
|
||||||
|
|
||||||
return this.prisma.agentProviderConfig.delete({
|
|
||||||
where: { id },
|
|
||||||
});
|
|
||||||
}
|
|
||||||
|
|
||||||
private toJsonValue(value: Record<string, unknown>): Prisma.InputJsonValue {
|
|
||||||
return value as Prisma.InputJsonValue;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -1,26 +0,0 @@
|
|||||||
import { IsBoolean, IsNotEmpty, IsObject, IsOptional, IsString, IsUUID } from "class-validator";
|
|
||||||
|
|
||||||
export class CreateAgentProviderDto {
|
|
||||||
@IsUUID()
|
|
||||||
workspaceId!: string;
|
|
||||||
|
|
||||||
@IsString()
|
|
||||||
@IsNotEmpty()
|
|
||||||
name!: string;
|
|
||||||
|
|
||||||
@IsString()
|
|
||||||
@IsNotEmpty()
|
|
||||||
provider!: string;
|
|
||||||
|
|
||||||
@IsString()
|
|
||||||
@IsNotEmpty()
|
|
||||||
gatewayUrl!: string;
|
|
||||||
|
|
||||||
@IsOptional()
|
|
||||||
@IsObject()
|
|
||||||
credentials?: Record<string, unknown>;
|
|
||||||
|
|
||||||
@IsOptional()
|
|
||||||
@IsBoolean()
|
|
||||||
isActive?: boolean;
|
|
||||||
}
|
|
||||||
@@ -1,30 +0,0 @@
|
|||||||
import { IsBoolean, IsNotEmpty, IsObject, IsOptional, IsString, IsUUID } from "class-validator";
|
|
||||||
|
|
||||||
export class UpdateAgentProviderDto {
|
|
||||||
@IsOptional()
|
|
||||||
@IsUUID()
|
|
||||||
workspaceId?: string;
|
|
||||||
|
|
||||||
@IsOptional()
|
|
||||||
@IsString()
|
|
||||||
@IsNotEmpty()
|
|
||||||
name?: string;
|
|
||||||
|
|
||||||
@IsOptional()
|
|
||||||
@IsString()
|
|
||||||
@IsNotEmpty()
|
|
||||||
provider?: string;
|
|
||||||
|
|
||||||
@IsOptional()
|
|
||||||
@IsString()
|
|
||||||
@IsNotEmpty()
|
|
||||||
gatewayUrl?: string;
|
|
||||||
|
|
||||||
@IsOptional()
|
|
||||||
@IsObject()
|
|
||||||
credentials?: Record<string, unknown>;
|
|
||||||
|
|
||||||
@IsOptional()
|
|
||||||
@IsBoolean()
|
|
||||||
isActive?: boolean;
|
|
||||||
}
|
|
||||||
@@ -1,172 +0,0 @@
|
|||||||
import { describe, it, expect, beforeEach, afterEach, vi } from "vitest";
|
|
||||||
import { AgentControlService } from "./agent-control.service";
|
|
||||||
import { PrismaService } from "../../prisma/prisma.service";
|
|
||||||
import { KillswitchService } from "../../killswitch/killswitch.service";
|
|
||||||
|
|
||||||
describe("AgentControlService", () => {
|
|
||||||
let service: AgentControlService;
|
|
||||||
let prisma: {
|
|
||||||
agentSessionTree: {
|
|
||||||
findUnique: ReturnType<typeof vi.fn>;
|
|
||||||
updateMany: ReturnType<typeof vi.fn>;
|
|
||||||
};
|
|
||||||
agentConversationMessage: {
|
|
||||||
create: ReturnType<typeof vi.fn>;
|
|
||||||
};
|
|
||||||
operatorAuditLog: {
|
|
||||||
create: ReturnType<typeof vi.fn>;
|
|
||||||
};
|
|
||||||
};
|
|
||||||
let killswitchService: {
|
|
||||||
killAgent: ReturnType<typeof vi.fn>;
|
|
||||||
};
|
|
||||||
|
|
||||||
beforeEach(() => {
|
|
||||||
prisma = {
|
|
||||||
agentSessionTree: {
|
|
||||||
findUnique: vi.fn(),
|
|
||||||
updateMany: vi.fn().mockResolvedValue({ count: 1 }),
|
|
||||||
},
|
|
||||||
agentConversationMessage: {
|
|
||||||
create: vi.fn().mockResolvedValue(undefined),
|
|
||||||
},
|
|
||||||
operatorAuditLog: {
|
|
||||||
create: vi.fn().mockResolvedValue(undefined),
|
|
||||||
},
|
|
||||||
};
|
|
||||||
|
|
||||||
killswitchService = {
|
|
||||||
killAgent: vi.fn().mockResolvedValue(undefined),
|
|
||||||
};
|
|
||||||
|
|
||||||
service = new AgentControlService(
|
|
||||||
prisma as unknown as PrismaService,
|
|
||||||
killswitchService as unknown as KillswitchService
|
|
||||||
);
|
|
||||||
});
|
|
||||||
|
|
||||||
afterEach(() => {
|
|
||||||
vi.clearAllMocks();
|
|
||||||
});
|
|
||||||
|
|
||||||
describe("injectMessage", () => {
|
|
||||||
it("creates conversation message and audit log when tree entry exists", async () => {
|
|
||||||
prisma.agentSessionTree.findUnique.mockResolvedValue({ id: "tree-1" });
|
|
||||||
|
|
||||||
await service.injectMessage("agent-123", "operator-abc", "Please continue");
|
|
||||||
|
|
||||||
expect(prisma.agentSessionTree.findUnique).toHaveBeenCalledWith({
|
|
||||||
where: { sessionId: "agent-123" },
|
|
||||||
select: { id: true },
|
|
||||||
});
|
|
||||||
expect(prisma.agentConversationMessage.create).toHaveBeenCalledWith({
|
|
||||||
data: {
|
|
||||||
sessionId: "agent-123",
|
|
||||||
role: "operator",
|
|
||||||
content: "Please continue",
|
|
||||||
provider: "internal",
|
|
||||||
metadata: {},
|
|
||||||
},
|
|
||||||
});
|
|
||||||
expect(prisma.operatorAuditLog.create).toHaveBeenCalledWith({
|
|
||||||
data: {
|
|
||||||
sessionId: "agent-123",
|
|
||||||
userId: "operator-abc",
|
|
||||||
provider: "internal",
|
|
||||||
action: "inject",
|
|
||||||
metadata: {
|
|
||||||
payload: {
|
|
||||||
message: "Please continue",
|
|
||||||
},
|
|
||||||
},
|
|
||||||
},
|
|
||||||
});
|
|
||||||
});
|
|
||||||
|
|
||||||
it("creates only audit log when no tree entry exists", async () => {
|
|
||||||
prisma.agentSessionTree.findUnique.mockResolvedValue(null);
|
|
||||||
|
|
||||||
await service.injectMessage("agent-456", "operator-def", "Nudge message");
|
|
||||||
|
|
||||||
expect(prisma.agentConversationMessage.create).not.toHaveBeenCalled();
|
|
||||||
expect(prisma.operatorAuditLog.create).toHaveBeenCalledWith({
|
|
||||||
data: {
|
|
||||||
sessionId: "agent-456",
|
|
||||||
userId: "operator-def",
|
|
||||||
provider: "internal",
|
|
||||||
action: "inject",
|
|
||||||
metadata: {
|
|
||||||
payload: {
|
|
||||||
message: "Nudge message",
|
|
||||||
},
|
|
||||||
},
|
|
||||||
},
|
|
||||||
});
|
|
||||||
});
|
|
||||||
});
|
|
||||||
|
|
||||||
describe("pauseAgent", () => {
|
|
||||||
it("updates tree status to paused and creates audit log", async () => {
|
|
||||||
await service.pauseAgent("agent-789", "operator-pause");
|
|
||||||
|
|
||||||
expect(prisma.agentSessionTree.updateMany).toHaveBeenCalledWith({
|
|
||||||
where: { sessionId: "agent-789" },
|
|
||||||
data: { status: "paused" },
|
|
||||||
});
|
|
||||||
expect(prisma.operatorAuditLog.create).toHaveBeenCalledWith({
|
|
||||||
data: {
|
|
||||||
sessionId: "agent-789",
|
|
||||||
userId: "operator-pause",
|
|
||||||
provider: "internal",
|
|
||||||
action: "pause",
|
|
||||||
metadata: {
|
|
||||||
payload: {},
|
|
||||||
},
|
|
||||||
},
|
|
||||||
});
|
|
||||||
});
|
|
||||||
});
|
|
||||||
|
|
||||||
describe("resumeAgent", () => {
|
|
||||||
it("updates tree status to running and creates audit log", async () => {
|
|
||||||
await service.resumeAgent("agent-321", "operator-resume");
|
|
||||||
|
|
||||||
expect(prisma.agentSessionTree.updateMany).toHaveBeenCalledWith({
|
|
||||||
where: { sessionId: "agent-321" },
|
|
||||||
data: { status: "running" },
|
|
||||||
});
|
|
||||||
expect(prisma.operatorAuditLog.create).toHaveBeenCalledWith({
|
|
||||||
data: {
|
|
||||||
sessionId: "agent-321",
|
|
||||||
userId: "operator-resume",
|
|
||||||
provider: "internal",
|
|
||||||
action: "resume",
|
|
||||||
metadata: {
|
|
||||||
payload: {},
|
|
||||||
},
|
|
||||||
},
|
|
||||||
});
|
|
||||||
});
|
|
||||||
});
|
|
||||||
|
|
||||||
describe("killAgent", () => {
|
|
||||||
it("delegates kill to killswitch and logs audit", async () => {
|
|
||||||
await service.killAgent("agent-654", "operator-kill", false);
|
|
||||||
|
|
||||||
expect(killswitchService.killAgent).toHaveBeenCalledWith("agent-654");
|
|
||||||
expect(prisma.operatorAuditLog.create).toHaveBeenCalledWith({
|
|
||||||
data: {
|
|
||||||
sessionId: "agent-654",
|
|
||||||
userId: "operator-kill",
|
|
||||||
provider: "internal",
|
|
||||||
action: "kill",
|
|
||||||
metadata: {
|
|
||||||
payload: {
|
|
||||||
force: false,
|
|
||||||
},
|
|
||||||
},
|
|
||||||
},
|
|
||||||
});
|
|
||||||
});
|
|
||||||
});
|
|
||||||
});
|
|
||||||
@@ -1,14 +1,10 @@
|
|||||||
import { Injectable } from "@nestjs/common";
|
import { Injectable } from "@nestjs/common";
|
||||||
import type { Prisma } from "@prisma/client";
|
import type { Prisma } from "@prisma/client";
|
||||||
import { KillswitchService } from "../../killswitch/killswitch.service";
|
|
||||||
import { PrismaService } from "../../prisma/prisma.service";
|
import { PrismaService } from "../../prisma/prisma.service";
|
||||||
|
|
||||||
@Injectable()
|
@Injectable()
|
||||||
export class AgentControlService {
|
export class AgentControlService {
|
||||||
constructor(
|
constructor(private readonly prisma: PrismaService) {}
|
||||||
private readonly prisma: PrismaService,
|
|
||||||
private readonly killswitchService: KillswitchService
|
|
||||||
) {}
|
|
||||||
|
|
||||||
private toJsonValue(value: Record<string, unknown>): Prisma.InputJsonValue {
|
private toJsonValue(value: Record<string, unknown>): Prisma.InputJsonValue {
|
||||||
return value as Prisma.InputJsonValue;
|
return value as Prisma.InputJsonValue;
|
||||||
@@ -17,7 +13,7 @@ export class AgentControlService {
|
|||||||
private async createOperatorAuditLog(
|
private async createOperatorAuditLog(
|
||||||
agentId: string,
|
agentId: string,
|
||||||
operatorId: string,
|
operatorId: string,
|
||||||
action: "inject" | "pause" | "resume" | "kill",
|
action: "inject" | "pause" | "resume",
|
||||||
payload: Record<string, unknown>
|
payload: Record<string, unknown>
|
||||||
): Promise<void> {
|
): Promise<void> {
|
||||||
await this.prisma.operatorAuditLog.create({
|
await this.prisma.operatorAuditLog.create({
|
||||||
@@ -69,9 +65,4 @@ export class AgentControlService {
|
|||||||
|
|
||||||
await this.createOperatorAuditLog(agentId, operatorId, "resume", {});
|
await this.createOperatorAuditLog(agentId, operatorId, "resume", {});
|
||||||
}
|
}
|
||||||
|
|
||||||
async killAgent(agentId: string, operatorId: string, force = true): Promise<void> {
|
|
||||||
await this.killswitchService.killAgent(agentId);
|
|
||||||
await this.createOperatorAuditLog(agentId, operatorId, "kill", { force });
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,103 +0,0 @@
|
|||||||
import { describe, it, expect, beforeEach, afterEach, vi } from "vitest";
|
|
||||||
import { AgentMessagesService } from "./agent-messages.service";
|
|
||||||
import { PrismaService } from "../../prisma/prisma.service";
|
|
||||||
|
|
||||||
describe("AgentMessagesService", () => {
|
|
||||||
let service: AgentMessagesService;
|
|
||||||
let prisma: {
|
|
||||||
agentConversationMessage: {
|
|
||||||
findMany: ReturnType<typeof vi.fn>;
|
|
||||||
count: ReturnType<typeof vi.fn>;
|
|
||||||
};
|
|
||||||
};
|
|
||||||
|
|
||||||
beforeEach(() => {
|
|
||||||
prisma = {
|
|
||||||
agentConversationMessage: {
|
|
||||||
findMany: vi.fn(),
|
|
||||||
count: vi.fn(),
|
|
||||||
},
|
|
||||||
};
|
|
||||||
|
|
||||||
service = new AgentMessagesService(prisma as unknown as PrismaService);
|
|
||||||
});
|
|
||||||
|
|
||||||
afterEach(() => {
|
|
||||||
vi.clearAllMocks();
|
|
||||||
});
|
|
||||||
|
|
||||||
describe("getMessages", () => {
|
|
||||||
it("returns paginated messages from Prisma", async () => {
|
|
||||||
const sessionId = "agent-123";
|
|
||||||
const messages = [
|
|
||||||
{
|
|
||||||
id: "msg-1",
|
|
||||||
sessionId,
|
|
||||||
provider: "internal",
|
|
||||||
role: "assistant",
|
|
||||||
content: "First message",
|
|
||||||
timestamp: new Date("2026-03-07T16:00:00.000Z"),
|
|
||||||
metadata: {},
|
|
||||||
},
|
|
||||||
{
|
|
||||||
id: "msg-2",
|
|
||||||
sessionId,
|
|
||||||
provider: "internal",
|
|
||||||
role: "user",
|
|
||||||
content: "Second message",
|
|
||||||
timestamp: new Date("2026-03-07T15:59:00.000Z"),
|
|
||||||
metadata: {},
|
|
||||||
},
|
|
||||||
];
|
|
||||||
|
|
||||||
prisma.agentConversationMessage.findMany.mockResolvedValue(messages);
|
|
||||||
prisma.agentConversationMessage.count.mockResolvedValue(2);
|
|
||||||
|
|
||||||
const result = await service.getMessages(sessionId, 50, 0);
|
|
||||||
|
|
||||||
expect(prisma.agentConversationMessage.findMany).toHaveBeenCalledWith({
|
|
||||||
where: { sessionId },
|
|
||||||
orderBy: { timestamp: "desc" },
|
|
||||||
take: 50,
|
|
||||||
skip: 0,
|
|
||||||
});
|
|
||||||
expect(prisma.agentConversationMessage.count).toHaveBeenCalledWith({ where: { sessionId } });
|
|
||||||
expect(result).toEqual({
|
|
||||||
messages,
|
|
||||||
total: 2,
|
|
||||||
});
|
|
||||||
});
|
|
||||||
|
|
||||||
it("applies limit and cursor (skip) correctly", async () => {
|
|
||||||
const sessionId = "agent-456";
|
|
||||||
const limit = 10;
|
|
||||||
const cursor = 20;
|
|
||||||
|
|
||||||
prisma.agentConversationMessage.findMany.mockResolvedValue([]);
|
|
||||||
prisma.agentConversationMessage.count.mockResolvedValue(42);
|
|
||||||
|
|
||||||
await service.getMessages(sessionId, limit, cursor);
|
|
||||||
|
|
||||||
expect(prisma.agentConversationMessage.findMany).toHaveBeenCalledWith({
|
|
||||||
where: { sessionId },
|
|
||||||
orderBy: { timestamp: "desc" },
|
|
||||||
take: limit,
|
|
||||||
skip: cursor,
|
|
||||||
});
|
|
||||||
});
|
|
||||||
|
|
||||||
it("returns empty messages array when no messages exist", async () => {
|
|
||||||
const sessionId = "agent-empty";
|
|
||||||
|
|
||||||
prisma.agentConversationMessage.findMany.mockResolvedValue([]);
|
|
||||||
prisma.agentConversationMessage.count.mockResolvedValue(0);
|
|
||||||
|
|
||||||
const result = await service.getMessages(sessionId, 25, 0);
|
|
||||||
|
|
||||||
expect(result).toEqual({
|
|
||||||
messages: [],
|
|
||||||
total: 0,
|
|
||||||
});
|
|
||||||
});
|
|
||||||
});
|
|
||||||
});
|
|
||||||
@@ -1,202 +0,0 @@
|
|||||||
import { Logger } from "@nestjs/common";
|
|
||||||
import type {
|
|
||||||
AgentMessage,
|
|
||||||
AgentSession,
|
|
||||||
AgentSessionList,
|
|
||||||
IAgentProvider,
|
|
||||||
InjectResult,
|
|
||||||
} from "@mosaic/shared";
|
|
||||||
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
|
|
||||||
import { AgentProviderRegistry } from "./agent-provider.registry";
|
|
||||||
import { InternalAgentProvider } from "./internal-agent.provider";
|
|
||||||
|
|
||||||
type MockProvider = IAgentProvider & {
|
|
||||||
listSessions: ReturnType<typeof vi.fn>;
|
|
||||||
getSession: ReturnType<typeof vi.fn>;
|
|
||||||
};
|
|
||||||
|
|
||||||
const emptyMessageStream = async function* (): AsyncIterable<AgentMessage> {
|
|
||||||
return;
|
|
||||||
};
|
|
||||||
|
|
||||||
const createProvider = (providerId: string, sessions: AgentSession[] = []): MockProvider => {
|
|
||||||
return {
|
|
||||||
providerId,
|
|
||||||
providerType: providerId,
|
|
||||||
displayName: providerId,
|
|
||||||
listSessions: vi.fn().mockResolvedValue({
|
|
||||||
sessions,
|
|
||||||
total: sessions.length,
|
|
||||||
} as AgentSessionList),
|
|
||||||
getSession: vi.fn().mockResolvedValue(null),
|
|
||||||
getMessages: vi.fn().mockResolvedValue([]),
|
|
||||||
injectMessage: vi.fn().mockResolvedValue({ accepted: true } as InjectResult),
|
|
||||||
pauseSession: vi.fn().mockResolvedValue(undefined),
|
|
||||||
resumeSession: vi.fn().mockResolvedValue(undefined),
|
|
||||||
killSession: vi.fn().mockResolvedValue(undefined),
|
|
||||||
streamMessages: vi.fn().mockReturnValue(emptyMessageStream()),
|
|
||||||
isAvailable: vi.fn().mockResolvedValue(true),
|
|
||||||
};
|
|
||||||
};
|
|
||||||
|
|
||||||
describe("AgentProviderRegistry", () => {
|
|
||||||
let registry: AgentProviderRegistry;
|
|
||||||
let internalProvider: MockProvider;
|
|
||||||
|
|
||||||
beforeEach(() => {
|
|
||||||
internalProvider = createProvider("internal");
|
|
||||||
registry = new AgentProviderRegistry(internalProvider as unknown as InternalAgentProvider);
|
|
||||||
});
|
|
||||||
|
|
||||||
afterEach(() => {
|
|
||||||
vi.restoreAllMocks();
|
|
||||||
});
|
|
||||||
|
|
||||||
it("registers InternalAgentProvider on module init", () => {
|
|
||||||
registry.onModuleInit();
|
|
||||||
|
|
||||||
expect(registry.getProvider("internal")).toBe(internalProvider);
|
|
||||||
});
|
|
||||||
|
|
||||||
it("registers providers and returns null for unknown provider ids", () => {
|
|
||||||
const externalProvider = createProvider("openclaw");
|
|
||||||
|
|
||||||
registry.registerProvider(externalProvider);
|
|
||||||
|
|
||||||
expect(registry.getProvider("openclaw")).toBe(externalProvider);
|
|
||||||
expect(registry.getProvider("missing")).toBeNull();
|
|
||||||
});
|
|
||||||
|
|
||||||
it("aggregates and sorts sessions from all providers", async () => {
|
|
||||||
const internalSessions: AgentSession[] = [
|
|
||||||
{
|
|
||||||
id: "session-older",
|
|
||||||
providerId: "internal",
|
|
||||||
providerType: "internal",
|
|
||||||
status: "active",
|
|
||||||
createdAt: new Date("2026-03-07T10:00:00.000Z"),
|
|
||||||
updatedAt: new Date("2026-03-07T10:10:00.000Z"),
|
|
||||||
},
|
|
||||||
];
|
|
||||||
|
|
||||||
const externalSessions: AgentSession[] = [
|
|
||||||
{
|
|
||||||
id: "session-newer",
|
|
||||||
providerId: "openclaw",
|
|
||||||
providerType: "external",
|
|
||||||
status: "paused",
|
|
||||||
createdAt: new Date("2026-03-07T09:00:00.000Z"),
|
|
||||||
updatedAt: new Date("2026-03-07T10:20:00.000Z"),
|
|
||||||
},
|
|
||||||
];
|
|
||||||
|
|
||||||
internalProvider.listSessions.mockResolvedValue({
|
|
||||||
sessions: internalSessions,
|
|
||||||
total: internalSessions.length,
|
|
||||||
} as AgentSessionList);
|
|
||||||
|
|
||||||
const externalProvider = createProvider("openclaw", externalSessions);
|
|
||||||
registry.onModuleInit();
|
|
||||||
registry.registerProvider(externalProvider);
|
|
||||||
|
|
||||||
const result = await registry.listAllSessions();
|
|
||||||
|
|
||||||
expect(result.map((session) => session.id)).toEqual(["session-newer", "session-older"]);
|
|
||||||
expect(internalProvider.listSessions).toHaveBeenCalledTimes(1);
|
|
||||||
expect(externalProvider.listSessions).toHaveBeenCalledTimes(1);
|
|
||||||
});
|
|
||||||
|
|
||||||
it("skips provider failures and logs warning", async () => {
|
|
||||||
const warnSpy = vi.spyOn(Logger.prototype, "warn").mockImplementation(() => undefined);
|
|
||||||
|
|
||||||
const healthyProvider = createProvider("healthy", [
|
|
||||||
{
|
|
||||||
id: "session-1",
|
|
||||||
providerId: "healthy",
|
|
||||||
providerType: "external",
|
|
||||||
status: "active",
|
|
||||||
createdAt: new Date("2026-03-07T11:00:00.000Z"),
|
|
||||||
updatedAt: new Date("2026-03-07T11:00:00.000Z"),
|
|
||||||
},
|
|
||||||
]);
|
|
||||||
|
|
||||||
const failingProvider = createProvider("failing");
|
|
||||||
failingProvider.listSessions.mockRejectedValue(new Error("provider offline"));
|
|
||||||
|
|
||||||
registry.onModuleInit();
|
|
||||||
registry.registerProvider(healthyProvider);
|
|
||||||
registry.registerProvider(failingProvider);
|
|
||||||
|
|
||||||
const result = await registry.listAllSessions();
|
|
||||||
|
|
||||||
expect(result).toHaveLength(1);
|
|
||||||
expect(result[0]?.id).toBe("session-1");
|
|
||||||
expect(warnSpy).toHaveBeenCalledWith(
|
|
||||||
expect.stringContaining("Failed to list sessions for provider failing")
|
|
||||||
);
|
|
||||||
});
|
|
||||||
|
|
||||||
it("finds a provider for an existing session", async () => {
|
|
||||||
const targetSession: AgentSession = {
|
|
||||||
id: "session-found",
|
|
||||||
providerId: "openclaw",
|
|
||||||
providerType: "external",
|
|
||||||
status: "active",
|
|
||||||
createdAt: new Date("2026-03-07T12:00:00.000Z"),
|
|
||||||
updatedAt: new Date("2026-03-07T12:10:00.000Z"),
|
|
||||||
};
|
|
||||||
|
|
||||||
const openclawProvider = createProvider("openclaw");
|
|
||||||
openclawProvider.getSession.mockResolvedValue(targetSession);
|
|
||||||
|
|
||||||
registry.onModuleInit();
|
|
||||||
registry.registerProvider(openclawProvider);
|
|
||||||
|
|
||||||
const result = await registry.getProviderForSession(targetSession.id);
|
|
||||||
|
|
||||||
expect(result).toEqual({
|
|
||||||
provider: openclawProvider,
|
|
||||||
session: targetSession,
|
|
||||||
});
|
|
||||||
expect(internalProvider.getSession).toHaveBeenCalledWith(targetSession.id);
|
|
||||||
expect(openclawProvider.getSession).toHaveBeenCalledWith(targetSession.id);
|
|
||||||
});
|
|
||||||
|
|
||||||
it("returns null when no provider has the requested session", async () => {
|
|
||||||
const openclawProvider = createProvider("openclaw");
|
|
||||||
|
|
||||||
registry.onModuleInit();
|
|
||||||
registry.registerProvider(openclawProvider);
|
|
||||||
|
|
||||||
await expect(registry.getProviderForSession("missing-session")).resolves.toBeNull();
|
|
||||||
});
|
|
||||||
|
|
||||||
it("continues searching providers when getSession throws", async () => {
|
|
||||||
const warnSpy = vi.spyOn(Logger.prototype, "warn").mockImplementation(() => undefined);
|
|
||||||
const failingProvider = createProvider("failing");
|
|
||||||
failingProvider.getSession.mockRejectedValue(new Error("provider timeout"));
|
|
||||||
|
|
||||||
const healthySession: AgentSession = {
|
|
||||||
id: "session-healthy",
|
|
||||||
providerId: "healthy",
|
|
||||||
providerType: "external",
|
|
||||||
status: "active",
|
|
||||||
createdAt: new Date("2026-03-07T12:15:00.000Z"),
|
|
||||||
updatedAt: new Date("2026-03-07T12:16:00.000Z"),
|
|
||||||
};
|
|
||||||
|
|
||||||
const healthyProvider = createProvider("healthy");
|
|
||||||
healthyProvider.getSession.mockResolvedValue(healthySession);
|
|
||||||
|
|
||||||
registry.onModuleInit();
|
|
||||||
registry.registerProvider(failingProvider);
|
|
||||||
registry.registerProvider(healthyProvider);
|
|
||||||
|
|
||||||
const result = await registry.getProviderForSession(healthySession.id);
|
|
||||||
|
|
||||||
expect(result).toEqual({ provider: healthyProvider, session: healthySession });
|
|
||||||
expect(warnSpy).toHaveBeenCalledWith(
|
|
||||||
expect.stringContaining("Failed to get session session-healthy for provider failing")
|
|
||||||
);
|
|
||||||
});
|
|
||||||
});
|
|
||||||
@@ -1,79 +0,0 @@
|
|||||||
import { Injectable, Logger, OnModuleInit } from "@nestjs/common";
|
|
||||||
import type { AgentSession, IAgentProvider } from "@mosaic/shared";
|
|
||||||
import { InternalAgentProvider } from "./internal-agent.provider";
|
|
||||||
|
|
||||||
@Injectable()
|
|
||||||
export class AgentProviderRegistry implements OnModuleInit {
|
|
||||||
private readonly logger = new Logger(AgentProviderRegistry.name);
|
|
||||||
private readonly providers = new Map<string, IAgentProvider>();
|
|
||||||
|
|
||||||
constructor(private readonly internalProvider: InternalAgentProvider) {}
|
|
||||||
|
|
||||||
onModuleInit(): void {
|
|
||||||
this.registerProvider(this.internalProvider);
|
|
||||||
}
|
|
||||||
|
|
||||||
registerProvider(provider: IAgentProvider): void {
|
|
||||||
const existingProvider = this.providers.get(provider.providerId);
|
|
||||||
if (existingProvider !== undefined) {
|
|
||||||
this.logger.warn(`Replacing existing provider registration for ${provider.providerId}`);
|
|
||||||
}
|
|
||||||
|
|
||||||
this.providers.set(provider.providerId, provider);
|
|
||||||
}
|
|
||||||
|
|
||||||
getProvider(providerId: string): IAgentProvider | null {
|
|
||||||
return this.providers.get(providerId) ?? null;
|
|
||||||
}
|
|
||||||
|
|
||||||
async getProviderForSession(
|
|
||||||
sessionId: string
|
|
||||||
): Promise<{ provider: IAgentProvider; session: AgentSession } | null> {
|
|
||||||
for (const provider of this.providers.values()) {
|
|
||||||
try {
|
|
||||||
const session = await provider.getSession(sessionId);
|
|
||||||
if (session !== null) {
|
|
||||||
return {
|
|
||||||
provider,
|
|
||||||
session,
|
|
||||||
};
|
|
||||||
}
|
|
||||||
} catch (error) {
|
|
||||||
this.logger.warn(
|
|
||||||
`Failed to get session ${sessionId} for provider ${provider.providerId}: ${this.toErrorMessage(error)}`
|
|
||||||
);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
return null;
|
|
||||||
}
|
|
||||||
|
|
||||||
async listAllSessions(): Promise<AgentSession[]> {
|
|
||||||
const providers = [...this.providers.values()];
|
|
||||||
const sessionsByProvider = await Promise.all(
|
|
||||||
providers.map(async (provider) => {
|
|
||||||
try {
|
|
||||||
const { sessions } = await provider.listSessions();
|
|
||||||
return sessions;
|
|
||||||
} catch (error) {
|
|
||||||
this.logger.warn(
|
|
||||||
`Failed to list sessions for provider ${provider.providerId}: ${this.toErrorMessage(error)}`
|
|
||||||
);
|
|
||||||
return [];
|
|
||||||
}
|
|
||||||
})
|
|
||||||
);
|
|
||||||
|
|
||||||
return sessionsByProvider
|
|
||||||
.flat()
|
|
||||||
.sort((left, right) => right.updatedAt.getTime() - left.updatedAt.getTime());
|
|
||||||
}
|
|
||||||
|
|
||||||
private toErrorMessage(error: unknown): string {
|
|
||||||
if (error instanceof Error) {
|
|
||||||
return error.message;
|
|
||||||
}
|
|
||||||
|
|
||||||
return String(error);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -1,245 +0,0 @@
|
|||||||
import { describe, it, expect, beforeEach, afterEach, vi } from "vitest";
|
|
||||||
import { AgentTreeService } from "./agent-tree.service";
|
|
||||||
import { PrismaService } from "../../prisma/prisma.service";
|
|
||||||
|
|
||||||
describe("AgentTreeService", () => {
|
|
||||||
let service: AgentTreeService;
|
|
||||||
let prisma: {
|
|
||||||
agentSessionTree: {
|
|
||||||
findMany: ReturnType<typeof vi.fn>;
|
|
||||||
count: ReturnType<typeof vi.fn>;
|
|
||||||
findUnique: ReturnType<typeof vi.fn>;
|
|
||||||
};
|
|
||||||
};
|
|
||||||
|
|
||||||
beforeEach(() => {
|
|
||||||
prisma = {
|
|
||||||
agentSessionTree: {
|
|
||||||
findMany: vi.fn(),
|
|
||||||
count: vi.fn(),
|
|
||||||
findUnique: vi.fn(),
|
|
||||||
},
|
|
||||||
};
|
|
||||||
|
|
||||||
service = new AgentTreeService(prisma as unknown as PrismaService);
|
|
||||||
});
|
|
||||||
|
|
||||||
afterEach(() => {
|
|
||||||
vi.clearAllMocks();
|
|
||||||
});
|
|
||||||
|
|
||||||
describe("listSessions", () => {
|
|
||||||
it("returns paginated sessions and cursor", async () => {
|
|
||||||
const sessions = [
|
|
||||||
{
|
|
||||||
id: "tree-2",
|
|
||||||
sessionId: "agent-2",
|
|
||||||
parentSessionId: null,
|
|
||||||
provider: "internal",
|
|
||||||
missionId: null,
|
|
||||||
taskId: "task-2",
|
|
||||||
taskSource: "queue",
|
|
||||||
agentType: "worker",
|
|
||||||
status: "running",
|
|
||||||
spawnedAt: new Date("2026-03-07T11:00:00.000Z"),
|
|
||||||
completedAt: null,
|
|
||||||
metadata: {},
|
|
||||||
},
|
|
||||||
{
|
|
||||||
id: "tree-1",
|
|
||||||
sessionId: "agent-1",
|
|
||||||
parentSessionId: null,
|
|
||||||
provider: "internal",
|
|
||||||
missionId: null,
|
|
||||||
taskId: "task-1",
|
|
||||||
taskSource: "queue",
|
|
||||||
agentType: "worker",
|
|
||||||
status: "running",
|
|
||||||
spawnedAt: new Date("2026-03-07T10:00:00.000Z"),
|
|
||||||
completedAt: null,
|
|
||||||
metadata: {},
|
|
||||||
},
|
|
||||||
];
|
|
||||||
|
|
||||||
prisma.agentSessionTree.findMany.mockResolvedValue(sessions);
|
|
||||||
prisma.agentSessionTree.count.mockResolvedValue(7);
|
|
||||||
|
|
||||||
const result = await service.listSessions(undefined, 2);
|
|
||||||
|
|
||||||
expect(prisma.agentSessionTree.findMany).toHaveBeenCalledWith({
|
|
||||||
where: undefined,
|
|
||||||
orderBy: [{ spawnedAt: "desc" }, { sessionId: "desc" }],
|
|
||||||
take: 2,
|
|
||||||
});
|
|
||||||
expect(prisma.agentSessionTree.count).toHaveBeenCalledWith();
|
|
||||||
expect(result.sessions).toEqual(sessions);
|
|
||||||
expect(result.total).toBe(7);
|
|
||||||
expect(result.cursor).toBeTypeOf("string");
|
|
||||||
});
|
|
||||||
|
|
||||||
it("applies cursor filter when provided", async () => {
|
|
||||||
prisma.agentSessionTree.findMany.mockResolvedValue([]);
|
|
||||||
prisma.agentSessionTree.count.mockResolvedValue(0);
|
|
||||||
|
|
||||||
const cursorDate = "2026-03-07T10:00:00.000Z";
|
|
||||||
const cursorSessionId = "agent-5";
|
|
||||||
const cursor = Buffer.from(
|
|
||||||
JSON.stringify({
|
|
||||||
spawnedAt: cursorDate,
|
|
||||||
sessionId: cursorSessionId,
|
|
||||||
}),
|
|
||||||
"utf8"
|
|
||||||
).toString("base64url");
|
|
||||||
|
|
||||||
await service.listSessions(cursor, 25);
|
|
||||||
|
|
||||||
expect(prisma.agentSessionTree.findMany).toHaveBeenCalledWith({
|
|
||||||
where: {
|
|
||||||
OR: [
|
|
||||||
{
|
|
||||||
spawnedAt: {
|
|
||||||
lt: new Date(cursorDate),
|
|
||||||
},
|
|
||||||
},
|
|
||||||
{
|
|
||||||
spawnedAt: new Date(cursorDate),
|
|
||||||
sessionId: {
|
|
||||||
lt: cursorSessionId,
|
|
||||||
},
|
|
||||||
},
|
|
||||||
],
|
|
||||||
},
|
|
||||||
orderBy: [{ spawnedAt: "desc" }, { sessionId: "desc" }],
|
|
||||||
take: 25,
|
|
||||||
});
|
|
||||||
});
|
|
||||||
|
|
||||||
it("ignores invalid cursor values", async () => {
|
|
||||||
prisma.agentSessionTree.findMany.mockResolvedValue([]);
|
|
||||||
prisma.agentSessionTree.count.mockResolvedValue(0);
|
|
||||||
|
|
||||||
await service.listSessions("invalid-cursor", 10);
|
|
||||||
|
|
||||||
expect(prisma.agentSessionTree.findMany).toHaveBeenCalledWith({
|
|
||||||
where: undefined,
|
|
||||||
orderBy: [{ spawnedAt: "desc" }, { sessionId: "desc" }],
|
|
||||||
take: 10,
|
|
||||||
});
|
|
||||||
});
|
|
||||||
});
|
|
||||||
|
|
||||||
describe("getSession", () => {
|
|
||||||
it("returns matching session entry", async () => {
|
|
||||||
const session = {
|
|
||||||
id: "tree-1",
|
|
||||||
sessionId: "agent-123",
|
|
||||||
parentSessionId: null,
|
|
||||||
provider: "internal",
|
|
||||||
missionId: null,
|
|
||||||
taskId: "task-1",
|
|
||||||
taskSource: "queue",
|
|
||||||
agentType: "worker",
|
|
||||||
status: "running",
|
|
||||||
spawnedAt: new Date("2026-03-07T11:00:00.000Z"),
|
|
||||||
completedAt: null,
|
|
||||||
metadata: {},
|
|
||||||
};
|
|
||||||
prisma.agentSessionTree.findUnique.mockResolvedValue(session);
|
|
||||||
|
|
||||||
const result = await service.getSession("agent-123");
|
|
||||||
|
|
||||||
expect(prisma.agentSessionTree.findUnique).toHaveBeenCalledWith({
|
|
||||||
where: { sessionId: "agent-123" },
|
|
||||||
});
|
|
||||||
expect(result).toEqual(session);
|
|
||||||
});
|
|
||||||
|
|
||||||
it("returns null when session does not exist", async () => {
|
|
||||||
prisma.agentSessionTree.findUnique.mockResolvedValue(null);
|
|
||||||
|
|
||||||
const result = await service.getSession("agent-missing");
|
|
||||||
|
|
||||||
expect(result).toBeNull();
|
|
||||||
});
|
|
||||||
});
|
|
||||||
|
|
||||||
describe("getTree", () => {
|
|
||||||
it("returns mapped entries from Prisma", async () => {
|
|
||||||
prisma.agentSessionTree.findMany.mockResolvedValue([
|
|
||||||
{
|
|
||||||
id: "tree-1",
|
|
||||||
sessionId: "agent-1",
|
|
||||||
parentSessionId: "agent-root",
|
|
||||||
provider: "internal",
|
|
||||||
missionId: "mission-1",
|
|
||||||
taskId: "task-1",
|
|
||||||
taskSource: "queue",
|
|
||||||
agentType: "worker",
|
|
||||||
status: "running",
|
|
||||||
spawnedAt: new Date("2026-03-07T10:00:00.000Z"),
|
|
||||||
completedAt: new Date("2026-03-07T11:00:00.000Z"),
|
|
||||||
metadata: {},
|
|
||||||
},
|
|
||||||
]);
|
|
||||||
|
|
||||||
const result = await service.getTree();
|
|
||||||
|
|
||||||
expect(prisma.agentSessionTree.findMany).toHaveBeenCalledWith({
|
|
||||||
orderBy: { spawnedAt: "desc" },
|
|
||||||
take: 200,
|
|
||||||
});
|
|
||||||
expect(result).toEqual([
|
|
||||||
{
|
|
||||||
sessionId: "agent-1",
|
|
||||||
parentSessionId: "agent-root",
|
|
||||||
status: "running",
|
|
||||||
agentType: "worker",
|
|
||||||
taskSource: "queue",
|
|
||||||
spawnedAt: "2026-03-07T10:00:00.000Z",
|
|
||||||
completedAt: "2026-03-07T11:00:00.000Z",
|
|
||||||
},
|
|
||||||
]);
|
|
||||||
});
|
|
||||||
|
|
||||||
it("returns empty array when no entries exist", async () => {
|
|
||||||
prisma.agentSessionTree.findMany.mockResolvedValue([]);
|
|
||||||
|
|
||||||
const result = await service.getTree();
|
|
||||||
|
|
||||||
expect(result).toEqual([]);
|
|
||||||
});
|
|
||||||
|
|
||||||
it("maps null parentSessionId and completedAt correctly", async () => {
|
|
||||||
prisma.agentSessionTree.findMany.mockResolvedValue([
|
|
||||||
{
|
|
||||||
id: "tree-2",
|
|
||||||
sessionId: "agent-root",
|
|
||||||
parentSessionId: null,
|
|
||||||
provider: "internal",
|
|
||||||
missionId: null,
|
|
||||||
taskId: null,
|
|
||||||
taskSource: null,
|
|
||||||
agentType: null,
|
|
||||||
status: "spawning",
|
|
||||||
spawnedAt: new Date("2026-03-07T09:00:00.000Z"),
|
|
||||||
completedAt: null,
|
|
||||||
metadata: {},
|
|
||||||
},
|
|
||||||
]);
|
|
||||||
|
|
||||||
const result = await service.getTree();
|
|
||||||
|
|
||||||
expect(result).toEqual([
|
|
||||||
{
|
|
||||||
sessionId: "agent-root",
|
|
||||||
parentSessionId: null,
|
|
||||||
status: "spawning",
|
|
||||||
agentType: null,
|
|
||||||
taskSource: null,
|
|
||||||
spawnedAt: "2026-03-07T09:00:00.000Z",
|
|
||||||
completedAt: null,
|
|
||||||
},
|
|
||||||
]);
|
|
||||||
});
|
|
||||||
});
|
|
||||||
});
|
|
||||||
@@ -1,146 +0,0 @@
|
|||||||
import { Injectable } from "@nestjs/common";
|
|
||||||
import type { AgentSessionTree, Prisma } from "@prisma/client";
|
|
||||||
import { AgentTreeResponseDto } from "./dto/agent-tree-response.dto";
|
|
||||||
import { PrismaService } from "../../prisma/prisma.service";
|
|
||||||
|
|
||||||
const DEFAULT_PAGE_LIMIT = 50;
|
|
||||||
const MAX_PAGE_LIMIT = 200;
|
|
||||||
|
|
||||||
interface SessionCursor {
|
|
||||||
spawnedAt: Date;
|
|
||||||
sessionId: string;
|
|
||||||
}
|
|
||||||
|
|
||||||
export interface AgentSessionTreeListResult {
|
|
||||||
sessions: AgentSessionTree[];
|
|
||||||
total: number;
|
|
||||||
cursor?: string;
|
|
||||||
}
|
|
||||||
|
|
||||||
@Injectable()
|
|
||||||
export class AgentTreeService {
|
|
||||||
constructor(private readonly prisma: PrismaService) {}
|
|
||||||
|
|
||||||
async listSessions(
|
|
||||||
cursor?: string,
|
|
||||||
limit = DEFAULT_PAGE_LIMIT
|
|
||||||
): Promise<AgentSessionTreeListResult> {
|
|
||||||
const safeLimit = this.normalizeLimit(limit);
|
|
||||||
const parsedCursor = this.parseCursor(cursor);
|
|
||||||
|
|
||||||
const where: Prisma.AgentSessionTreeWhereInput | undefined = parsedCursor
|
|
||||||
? {
|
|
||||||
OR: [
|
|
||||||
{
|
|
||||||
spawnedAt: {
|
|
||||||
lt: parsedCursor.spawnedAt,
|
|
||||||
},
|
|
||||||
},
|
|
||||||
{
|
|
||||||
spawnedAt: parsedCursor.spawnedAt,
|
|
||||||
sessionId: {
|
|
||||||
lt: parsedCursor.sessionId,
|
|
||||||
},
|
|
||||||
},
|
|
||||||
],
|
|
||||||
}
|
|
||||||
: undefined;
|
|
||||||
|
|
||||||
const [sessions, total] = await Promise.all([
|
|
||||||
this.prisma.agentSessionTree.findMany({
|
|
||||||
where,
|
|
||||||
orderBy: [{ spawnedAt: "desc" }, { sessionId: "desc" }],
|
|
||||||
take: safeLimit,
|
|
||||||
}),
|
|
||||||
this.prisma.agentSessionTree.count(),
|
|
||||||
]);
|
|
||||||
|
|
||||||
const nextCursor =
|
|
||||||
sessions.length === safeLimit
|
|
||||||
? this.serializeCursor(sessions[sessions.length - 1])
|
|
||||||
: undefined;
|
|
||||||
|
|
||||||
return {
|
|
||||||
sessions,
|
|
||||||
total,
|
|
||||||
...(nextCursor !== undefined ? { cursor: nextCursor } : {}),
|
|
||||||
};
|
|
||||||
}
|
|
||||||
|
|
||||||
async getSession(sessionId: string): Promise<AgentSessionTree | null> {
|
|
||||||
return this.prisma.agentSessionTree.findUnique({
|
|
||||||
where: { sessionId },
|
|
||||||
});
|
|
||||||
}
|
|
||||||
|
|
||||||
async getTree(): Promise<AgentTreeResponseDto[]> {
|
|
||||||
const entries = await this.prisma.agentSessionTree.findMany({
|
|
||||||
orderBy: { spawnedAt: "desc" },
|
|
||||||
take: 200,
|
|
||||||
});
|
|
||||||
|
|
||||||
const response: AgentTreeResponseDto[] = [];
|
|
||||||
for (const entry of entries) {
|
|
||||||
response.push({
|
|
||||||
sessionId: entry.sessionId,
|
|
||||||
parentSessionId: entry.parentSessionId ?? null,
|
|
||||||
status: entry.status,
|
|
||||||
agentType: entry.agentType ?? null,
|
|
||||||
taskSource: entry.taskSource ?? null,
|
|
||||||
spawnedAt: entry.spawnedAt.toISOString(),
|
|
||||||
completedAt: entry.completedAt?.toISOString() ?? null,
|
|
||||||
});
|
|
||||||
}
|
|
||||||
|
|
||||||
return response;
|
|
||||||
}
|
|
||||||
|
|
||||||
private normalizeLimit(limit: number): number {
|
|
||||||
const normalized = Number.isFinite(limit) ? Math.trunc(limit) : DEFAULT_PAGE_LIMIT;
|
|
||||||
if (normalized < 1) {
|
|
||||||
return 1;
|
|
||||||
}
|
|
||||||
|
|
||||||
return Math.min(normalized, MAX_PAGE_LIMIT);
|
|
||||||
}
|
|
||||||
|
|
||||||
private serializeCursor(entry: Pick<AgentSessionTree, "spawnedAt" | "sessionId">): string {
|
|
||||||
return Buffer.from(
|
|
||||||
JSON.stringify({
|
|
||||||
spawnedAt: entry.spawnedAt.toISOString(),
|
|
||||||
sessionId: entry.sessionId,
|
|
||||||
}),
|
|
||||||
"utf8"
|
|
||||||
).toString("base64url");
|
|
||||||
}
|
|
||||||
|
|
||||||
private parseCursor(cursor?: string): SessionCursor | null {
|
|
||||||
if (!cursor) {
|
|
||||||
return null;
|
|
||||||
}
|
|
||||||
|
|
||||||
try {
|
|
||||||
const decoded = Buffer.from(cursor, "base64url").toString("utf8");
|
|
||||||
const parsed = JSON.parse(decoded) as {
|
|
||||||
spawnedAt?: string;
|
|
||||||
sessionId?: string;
|
|
||||||
};
|
|
||||||
|
|
||||||
if (typeof parsed.spawnedAt !== "string" || typeof parsed.sessionId !== "string") {
|
|
||||||
return null;
|
|
||||||
}
|
|
||||||
|
|
||||||
const spawnedAt = new Date(parsed.spawnedAt);
|
|
||||||
if (Number.isNaN(spawnedAt.getTime())) {
|
|
||||||
return null;
|
|
||||||
}
|
|
||||||
|
|
||||||
return {
|
|
||||||
spawnedAt,
|
|
||||||
sessionId: parsed.sessionId,
|
|
||||||
};
|
|
||||||
} catch {
|
|
||||||
return null;
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -7,7 +7,6 @@ import { KillswitchService } from "../../killswitch/killswitch.service";
|
|||||||
import { AgentEventsService } from "./agent-events.service";
|
import { AgentEventsService } from "./agent-events.service";
|
||||||
import { AgentMessagesService } from "./agent-messages.service";
|
import { AgentMessagesService } from "./agent-messages.service";
|
||||||
import { AgentControlService } from "./agent-control.service";
|
import { AgentControlService } from "./agent-control.service";
|
||||||
import { AgentTreeService } from "./agent-tree.service";
|
|
||||||
import type { KillAllResult } from "../../killswitch/killswitch.service";
|
import type { KillAllResult } from "../../killswitch/killswitch.service";
|
||||||
|
|
||||||
describe("AgentsController - Killswitch Endpoints", () => {
|
describe("AgentsController - Killswitch Endpoints", () => {
|
||||||
@@ -42,9 +41,6 @@ describe("AgentsController - Killswitch Endpoints", () => {
|
|||||||
pauseAgent: ReturnType<typeof vi.fn>;
|
pauseAgent: ReturnType<typeof vi.fn>;
|
||||||
resumeAgent: ReturnType<typeof vi.fn>;
|
resumeAgent: ReturnType<typeof vi.fn>;
|
||||||
};
|
};
|
||||||
let mockTreeService: {
|
|
||||||
getTree: ReturnType<typeof vi.fn>;
|
|
||||||
};
|
|
||||||
|
|
||||||
beforeEach(() => {
|
beforeEach(() => {
|
||||||
mockKillswitchService = {
|
mockKillswitchService = {
|
||||||
@@ -93,10 +89,6 @@ describe("AgentsController - Killswitch Endpoints", () => {
|
|||||||
resumeAgent: vi.fn().mockResolvedValue(undefined),
|
resumeAgent: vi.fn().mockResolvedValue(undefined),
|
||||||
};
|
};
|
||||||
|
|
||||||
mockTreeService = {
|
|
||||||
getTree: vi.fn().mockResolvedValue([]),
|
|
||||||
};
|
|
||||||
|
|
||||||
controller = new AgentsController(
|
controller = new AgentsController(
|
||||||
mockQueueService as unknown as QueueService,
|
mockQueueService as unknown as QueueService,
|
||||||
mockSpawnerService as unknown as AgentSpawnerService,
|
mockSpawnerService as unknown as AgentSpawnerService,
|
||||||
@@ -104,8 +96,7 @@ describe("AgentsController - Killswitch Endpoints", () => {
|
|||||||
mockKillswitchService as unknown as KillswitchService,
|
mockKillswitchService as unknown as KillswitchService,
|
||||||
mockEventsService as unknown as AgentEventsService,
|
mockEventsService as unknown as AgentEventsService,
|
||||||
mockMessagesService as unknown as AgentMessagesService,
|
mockMessagesService as unknown as AgentMessagesService,
|
||||||
mockControlService as unknown as AgentControlService,
|
mockControlService as unknown as AgentControlService
|
||||||
mockTreeService as unknown as AgentTreeService
|
|
||||||
);
|
);
|
||||||
});
|
});
|
||||||
|
|
||||||
|
|||||||
@@ -6,7 +6,6 @@ import { KillswitchService } from "../../killswitch/killswitch.service";
|
|||||||
import { AgentEventsService } from "./agent-events.service";
|
import { AgentEventsService } from "./agent-events.service";
|
||||||
import { AgentMessagesService } from "./agent-messages.service";
|
import { AgentMessagesService } from "./agent-messages.service";
|
||||||
import { AgentControlService } from "./agent-control.service";
|
import { AgentControlService } from "./agent-control.service";
|
||||||
import { AgentTreeService } from "./agent-tree.service";
|
|
||||||
import { describe, it, expect, beforeEach, afterEach, vi } from "vitest";
|
import { describe, it, expect, beforeEach, afterEach, vi } from "vitest";
|
||||||
|
|
||||||
describe("AgentsController", () => {
|
describe("AgentsController", () => {
|
||||||
@@ -43,9 +42,6 @@ describe("AgentsController", () => {
|
|||||||
pauseAgent: ReturnType<typeof vi.fn>;
|
pauseAgent: ReturnType<typeof vi.fn>;
|
||||||
resumeAgent: ReturnType<typeof vi.fn>;
|
resumeAgent: ReturnType<typeof vi.fn>;
|
||||||
};
|
};
|
||||||
let treeService: {
|
|
||||||
getTree: ReturnType<typeof vi.fn>;
|
|
||||||
};
|
|
||||||
|
|
||||||
beforeEach(() => {
|
beforeEach(() => {
|
||||||
// Create mock services
|
// Create mock services
|
||||||
@@ -97,10 +93,6 @@ describe("AgentsController", () => {
|
|||||||
resumeAgent: vi.fn().mockResolvedValue(undefined),
|
resumeAgent: vi.fn().mockResolvedValue(undefined),
|
||||||
};
|
};
|
||||||
|
|
||||||
treeService = {
|
|
||||||
getTree: vi.fn().mockResolvedValue([]),
|
|
||||||
};
|
|
||||||
|
|
||||||
// Create controller with mocked services
|
// Create controller with mocked services
|
||||||
controller = new AgentsController(
|
controller = new AgentsController(
|
||||||
queueService as unknown as QueueService,
|
queueService as unknown as QueueService,
|
||||||
@@ -109,8 +101,7 @@ describe("AgentsController", () => {
|
|||||||
killswitchService as unknown as KillswitchService,
|
killswitchService as unknown as KillswitchService,
|
||||||
eventsService as unknown as AgentEventsService,
|
eventsService as unknown as AgentEventsService,
|
||||||
messagesService as unknown as AgentMessagesService,
|
messagesService as unknown as AgentMessagesService,
|
||||||
controlService as unknown as AgentControlService,
|
controlService as unknown as AgentControlService
|
||||||
treeService as unknown as AgentTreeService
|
|
||||||
);
|
);
|
||||||
});
|
});
|
||||||
|
|
||||||
@@ -122,27 +113,6 @@ describe("AgentsController", () => {
|
|||||||
expect(controller).toBeDefined();
|
expect(controller).toBeDefined();
|
||||||
});
|
});
|
||||||
|
|
||||||
describe("getAgentTree", () => {
|
|
||||||
it("should return tree entries", async () => {
|
|
||||||
const entries = [
|
|
||||||
{
|
|
||||||
sessionId: "agent-1",
|
|
||||||
parentSessionId: null,
|
|
||||||
status: "running",
|
|
||||||
agentType: "worker",
|
|
||||||
taskSource: "internal",
|
|
||||||
spawnedAt: "2026-03-07T00:00:00.000Z",
|
|
||||||
completedAt: null,
|
|
||||||
},
|
|
||||||
];
|
|
||||||
|
|
||||||
treeService.getTree.mockResolvedValue(entries);
|
|
||||||
|
|
||||||
await expect(controller.getAgentTree()).resolves.toEqual(entries);
|
|
||||||
expect(treeService.getTree).toHaveBeenCalledTimes(1);
|
|
||||||
});
|
|
||||||
});
|
|
||||||
|
|
||||||
describe("listAgents", () => {
|
describe("listAgents", () => {
|
||||||
it("should return empty array when no agents exist", () => {
|
it("should return empty array when no agents exist", () => {
|
||||||
// Arrange
|
// Arrange
|
||||||
|
|||||||
@@ -30,8 +30,6 @@ import { AgentEventsService } from "./agent-events.service";
|
|||||||
import { GetMessagesQueryDto } from "./dto/get-messages-query.dto";
|
import { GetMessagesQueryDto } from "./dto/get-messages-query.dto";
|
||||||
import { AgentMessagesService } from "./agent-messages.service";
|
import { AgentMessagesService } from "./agent-messages.service";
|
||||||
import { AgentControlService } from "./agent-control.service";
|
import { AgentControlService } from "./agent-control.service";
|
||||||
import { AgentTreeService } from "./agent-tree.service";
|
|
||||||
import { AgentTreeResponseDto } from "./dto/agent-tree-response.dto";
|
|
||||||
import { InjectAgentDto } from "./dto/inject-agent.dto";
|
import { InjectAgentDto } from "./dto/inject-agent.dto";
|
||||||
import { PauseAgentDto, ResumeAgentDto } from "./dto/control-agent.dto";
|
import { PauseAgentDto, ResumeAgentDto } from "./dto/control-agent.dto";
|
||||||
|
|
||||||
@@ -58,8 +56,7 @@ export class AgentsController {
|
|||||||
private readonly killswitchService: KillswitchService,
|
private readonly killswitchService: KillswitchService,
|
||||||
private readonly eventsService: AgentEventsService,
|
private readonly eventsService: AgentEventsService,
|
||||||
private readonly messagesService: AgentMessagesService,
|
private readonly messagesService: AgentMessagesService,
|
||||||
private readonly agentControlService: AgentControlService,
|
private readonly agentControlService: AgentControlService
|
||||||
private readonly agentTreeService: AgentTreeService
|
|
||||||
) {}
|
) {}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -81,7 +78,6 @@ export class AgentsController {
|
|||||||
// Spawn agent using spawner service
|
// Spawn agent using spawner service
|
||||||
const spawnResponse = this.spawnerService.spawnAgent({
|
const spawnResponse = this.spawnerService.spawnAgent({
|
||||||
taskId: dto.taskId,
|
taskId: dto.taskId,
|
||||||
...(dto.parentAgentId !== undefined ? { parentAgentId: dto.parentAgentId } : {}),
|
|
||||||
agentType: dto.agentType,
|
agentType: dto.agentType,
|
||||||
context: dto.context,
|
context: dto.context,
|
||||||
});
|
});
|
||||||
@@ -156,13 +152,6 @@ export class AgentsController {
|
|||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
@Get("tree")
|
|
||||||
@UseGuards(OrchestratorApiKeyGuard)
|
|
||||||
@Throttle({ default: { limit: 200, ttl: 60000 } })
|
|
||||||
async getAgentTree(): Promise<AgentTreeResponseDto[]> {
|
|
||||||
return this.agentTreeService.getTree();
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* List all agents
|
* List all agents
|
||||||
* @returns Array of all agent sessions with their status
|
* @returns Array of all agent sessions with their status
|
||||||
|
|||||||
@@ -9,9 +9,6 @@ import { AgentEventsService } from "./agent-events.service";
|
|||||||
import { PrismaModule } from "../../prisma/prisma.module";
|
import { PrismaModule } from "../../prisma/prisma.module";
|
||||||
import { AgentMessagesService } from "./agent-messages.service";
|
import { AgentMessagesService } from "./agent-messages.service";
|
||||||
import { AgentControlService } from "./agent-control.service";
|
import { AgentControlService } from "./agent-control.service";
|
||||||
import { AgentTreeService } from "./agent-tree.service";
|
|
||||||
import { InternalAgentProvider } from "./internal-agent.provider";
|
|
||||||
import { AgentProviderRegistry } from "./agent-provider.registry";
|
|
||||||
|
|
||||||
@Module({
|
@Module({
|
||||||
imports: [QueueModule, SpawnerModule, KillswitchModule, ValkeyModule, PrismaModule],
|
imports: [QueueModule, SpawnerModule, KillswitchModule, ValkeyModule, PrismaModule],
|
||||||
@@ -21,10 +18,6 @@ import { AgentProviderRegistry } from "./agent-provider.registry";
|
|||||||
AgentEventsService,
|
AgentEventsService,
|
||||||
AgentMessagesService,
|
AgentMessagesService,
|
||||||
AgentControlService,
|
AgentControlService,
|
||||||
AgentTreeService,
|
|
||||||
InternalAgentProvider,
|
|
||||||
AgentProviderRegistry,
|
|
||||||
],
|
],
|
||||||
exports: [InternalAgentProvider, AgentProviderRegistry],
|
|
||||||
})
|
})
|
||||||
export class AgentsModule {}
|
export class AgentsModule {}
|
||||||
|
|||||||
@@ -1,9 +0,0 @@
|
|||||||
export class AgentTreeResponseDto {
|
|
||||||
sessionId!: string;
|
|
||||||
parentSessionId!: string | null;
|
|
||||||
status!: string;
|
|
||||||
agentType!: string | null;
|
|
||||||
taskSource!: string | null;
|
|
||||||
spawnedAt!: string;
|
|
||||||
completedAt!: string | null;
|
|
||||||
}
|
|
||||||
@@ -116,10 +116,6 @@ export class SpawnAgentDto {
|
|||||||
@IsOptional()
|
@IsOptional()
|
||||||
@IsIn(["strict", "standard", "minimal", "custom"])
|
@IsIn(["strict", "standard", "minimal", "custom"])
|
||||||
gateProfile?: GateProfileType;
|
gateProfile?: GateProfileType;
|
||||||
|
|
||||||
@IsOptional()
|
|
||||||
@IsString()
|
|
||||||
parentAgentId?: string;
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|||||||
@@ -1,216 +0,0 @@
|
|||||||
import { beforeEach, describe, expect, it, vi } from "vitest";
|
|
||||||
import type { AgentConversationMessage, AgentSessionTree } from "@prisma/client";
|
|
||||||
import { AgentControlService } from "./agent-control.service";
|
|
||||||
import { AgentMessagesService } from "./agent-messages.service";
|
|
||||||
import { AgentTreeService } from "./agent-tree.service";
|
|
||||||
import { InternalAgentProvider } from "./internal-agent.provider";
|
|
||||||
|
|
||||||
describe("InternalAgentProvider", () => {
|
|
||||||
let provider: InternalAgentProvider;
|
|
||||||
let messagesService: {
|
|
||||||
getMessages: ReturnType<typeof vi.fn>;
|
|
||||||
getReplayMessages: ReturnType<typeof vi.fn>;
|
|
||||||
getMessagesAfter: ReturnType<typeof vi.fn>;
|
|
||||||
};
|
|
||||||
let controlService: {
|
|
||||||
injectMessage: ReturnType<typeof vi.fn>;
|
|
||||||
pauseAgent: ReturnType<typeof vi.fn>;
|
|
||||||
resumeAgent: ReturnType<typeof vi.fn>;
|
|
||||||
killAgent: ReturnType<typeof vi.fn>;
|
|
||||||
};
|
|
||||||
let treeService: {
|
|
||||||
listSessions: ReturnType<typeof vi.fn>;
|
|
||||||
getSession: ReturnType<typeof vi.fn>;
|
|
||||||
};
|
|
||||||
|
|
||||||
beforeEach(() => {
|
|
||||||
messagesService = {
|
|
||||||
getMessages: vi.fn(),
|
|
||||||
getReplayMessages: vi.fn(),
|
|
||||||
getMessagesAfter: vi.fn(),
|
|
||||||
};
|
|
||||||
|
|
||||||
controlService = {
|
|
||||||
injectMessage: vi.fn().mockResolvedValue(undefined),
|
|
||||||
pauseAgent: vi.fn().mockResolvedValue(undefined),
|
|
||||||
resumeAgent: vi.fn().mockResolvedValue(undefined),
|
|
||||||
killAgent: vi.fn().mockResolvedValue(undefined),
|
|
||||||
};
|
|
||||||
|
|
||||||
treeService = {
|
|
||||||
listSessions: vi.fn(),
|
|
||||||
getSession: vi.fn(),
|
|
||||||
};
|
|
||||||
|
|
||||||
provider = new InternalAgentProvider(
|
|
||||||
messagesService as unknown as AgentMessagesService,
|
|
||||||
controlService as unknown as AgentControlService,
|
|
||||||
treeService as unknown as AgentTreeService
|
|
||||||
);
|
|
||||||
});
|
|
||||||
|
|
||||||
it("maps paginated sessions", async () => {
|
|
||||||
const sessionEntry: AgentSessionTree = {
|
|
||||||
id: "tree-1",
|
|
||||||
sessionId: "session-1",
|
|
||||||
parentSessionId: "parent-1",
|
|
||||||
provider: "internal",
|
|
||||||
missionId: null,
|
|
||||||
taskId: "task-123",
|
|
||||||
taskSource: "queue",
|
|
||||||
agentType: "worker",
|
|
||||||
status: "running",
|
|
||||||
spawnedAt: new Date("2026-03-07T10:00:00.000Z"),
|
|
||||||
completedAt: null,
|
|
||||||
metadata: { branch: "feat/test" },
|
|
||||||
};
|
|
||||||
|
|
||||||
treeService.listSessions.mockResolvedValue({
|
|
||||||
sessions: [sessionEntry],
|
|
||||||
total: 1,
|
|
||||||
cursor: "next-cursor",
|
|
||||||
});
|
|
||||||
|
|
||||||
const result = await provider.listSessions("cursor-1", 25);
|
|
||||||
|
|
||||||
expect(treeService.listSessions).toHaveBeenCalledWith("cursor-1", 25);
|
|
||||||
expect(result).toEqual({
|
|
||||||
sessions: [
|
|
||||||
{
|
|
||||||
id: "session-1",
|
|
||||||
providerId: "internal",
|
|
||||||
providerType: "internal",
|
|
||||||
label: "task-123",
|
|
||||||
status: "active",
|
|
||||||
parentSessionId: "parent-1",
|
|
||||||
createdAt: new Date("2026-03-07T10:00:00.000Z"),
|
|
||||||
updatedAt: new Date("2026-03-07T10:00:00.000Z"),
|
|
||||||
metadata: { branch: "feat/test" },
|
|
||||||
},
|
|
||||||
],
|
|
||||||
total: 1,
|
|
||||||
cursor: "next-cursor",
|
|
||||||
});
|
|
||||||
});
|
|
||||||
|
|
||||||
it("returns null for missing session", async () => {
|
|
||||||
treeService.getSession.mockResolvedValue(null);
|
|
||||||
|
|
||||||
const result = await provider.getSession("missing-session");
|
|
||||||
|
|
||||||
expect(treeService.getSession).toHaveBeenCalledWith("missing-session");
|
|
||||||
expect(result).toBeNull();
|
|
||||||
});
|
|
||||||
|
|
||||||
it("maps message history and parses skip cursor", async () => {
|
|
||||||
const message: AgentConversationMessage = {
|
|
||||||
id: "msg-1",
|
|
||||||
sessionId: "session-1",
|
|
||||||
provider: "internal",
|
|
||||||
role: "agent",
|
|
||||||
content: "hello",
|
|
||||||
timestamp: new Date("2026-03-07T10:05:00.000Z"),
|
|
||||||
metadata: { tokens: 42 },
|
|
||||||
};
|
|
||||||
|
|
||||||
messagesService.getMessages.mockResolvedValue({
|
|
||||||
messages: [message],
|
|
||||||
total: 10,
|
|
||||||
});
|
|
||||||
|
|
||||||
const result = await provider.getMessages("session-1", 30, "2");
|
|
||||||
|
|
||||||
expect(messagesService.getMessages).toHaveBeenCalledWith("session-1", 30, 2);
|
|
||||||
expect(result).toEqual([
|
|
||||||
{
|
|
||||||
id: "msg-1",
|
|
||||||
sessionId: "session-1",
|
|
||||||
role: "assistant",
|
|
||||||
content: "hello",
|
|
||||||
timestamp: new Date("2026-03-07T10:05:00.000Z"),
|
|
||||||
metadata: { tokens: 42 },
|
|
||||||
},
|
|
||||||
]);
|
|
||||||
});
|
|
||||||
|
|
||||||
it("routes control operations through AgentControlService", async () => {
|
|
||||||
const injectResult = await provider.injectMessage("session-1", "new instruction");
|
|
||||||
|
|
||||||
await provider.pauseSession("session-1");
|
|
||||||
await provider.resumeSession("session-1");
|
|
||||||
await provider.killSession("session-1", false);
|
|
||||||
|
|
||||||
expect(controlService.injectMessage).toHaveBeenCalledWith(
|
|
||||||
"session-1",
|
|
||||||
"internal-provider",
|
|
||||||
"new instruction"
|
|
||||||
);
|
|
||||||
expect(injectResult).toEqual({ accepted: true });
|
|
||||||
expect(controlService.pauseAgent).toHaveBeenCalledWith("session-1", "internal-provider");
|
|
||||||
expect(controlService.resumeAgent).toHaveBeenCalledWith("session-1", "internal-provider");
|
|
||||||
expect(controlService.killAgent).toHaveBeenCalledWith("session-1", "internal-provider", false);
|
|
||||||
});
|
|
||||||
|
|
||||||
it("streams replay and incremental messages", async () => {
|
|
||||||
const replayMessage: AgentConversationMessage = {
|
|
||||||
id: "msg-replay",
|
|
||||||
sessionId: "session-1",
|
|
||||||
provider: "internal",
|
|
||||||
role: "agent",
|
|
||||||
content: "replay",
|
|
||||||
timestamp: new Date("2026-03-07T10:00:00.000Z"),
|
|
||||||
metadata: {},
|
|
||||||
};
|
|
||||||
const incrementalMessage: AgentConversationMessage = {
|
|
||||||
id: "msg-live",
|
|
||||||
sessionId: "session-1",
|
|
||||||
provider: "internal",
|
|
||||||
role: "operator",
|
|
||||||
content: "live",
|
|
||||||
timestamp: new Date("2026-03-07T10:00:01.000Z"),
|
|
||||||
metadata: {},
|
|
||||||
};
|
|
||||||
|
|
||||||
messagesService.getReplayMessages.mockResolvedValue([replayMessage]);
|
|
||||||
messagesService.getMessagesAfter
|
|
||||||
.mockResolvedValueOnce([incrementalMessage])
|
|
||||||
.mockResolvedValueOnce([]);
|
|
||||||
|
|
||||||
const iterator = provider.streamMessages("session-1")[Symbol.asyncIterator]();
|
|
||||||
|
|
||||||
const first = await iterator.next();
|
|
||||||
const second = await iterator.next();
|
|
||||||
|
|
||||||
expect(first.done).toBe(false);
|
|
||||||
expect(first.value).toEqual({
|
|
||||||
id: "msg-replay",
|
|
||||||
sessionId: "session-1",
|
|
||||||
role: "assistant",
|
|
||||||
content: "replay",
|
|
||||||
timestamp: new Date("2026-03-07T10:00:00.000Z"),
|
|
||||||
metadata: {},
|
|
||||||
});
|
|
||||||
expect(second.done).toBe(false);
|
|
||||||
expect(second.value).toEqual({
|
|
||||||
id: "msg-live",
|
|
||||||
sessionId: "session-1",
|
|
||||||
role: "user",
|
|
||||||
content: "live",
|
|
||||||
timestamp: new Date("2026-03-07T10:00:01.000Z"),
|
|
||||||
metadata: {},
|
|
||||||
});
|
|
||||||
|
|
||||||
await iterator.return?.();
|
|
||||||
|
|
||||||
expect(messagesService.getReplayMessages).toHaveBeenCalledWith("session-1", 50);
|
|
||||||
expect(messagesService.getMessagesAfter).toHaveBeenCalledWith(
|
|
||||||
"session-1",
|
|
||||||
new Date("2026-03-07T10:00:00.000Z"),
|
|
||||||
"msg-replay"
|
|
||||||
);
|
|
||||||
});
|
|
||||||
|
|
||||||
it("reports provider availability", async () => {
|
|
||||||
await expect(provider.isAvailable()).resolves.toBe(true);
|
|
||||||
});
|
|
||||||
});
|
|
||||||
@@ -1,218 +0,0 @@
|
|||||||
import { Injectable } from "@nestjs/common";
|
|
||||||
import type {
|
|
||||||
AgentMessage,
|
|
||||||
AgentMessageRole,
|
|
||||||
AgentSession,
|
|
||||||
AgentSessionList,
|
|
||||||
AgentSessionStatus,
|
|
||||||
IAgentProvider,
|
|
||||||
InjectResult,
|
|
||||||
} from "@mosaic/shared";
|
|
||||||
import type { AgentConversationMessage, AgentSessionTree } from "@prisma/client";
|
|
||||||
import { AgentControlService } from "./agent-control.service";
|
|
||||||
import { AgentMessagesService } from "./agent-messages.service";
|
|
||||||
import { AgentTreeService } from "./agent-tree.service";
|
|
||||||
|
|
||||||
const DEFAULT_SESSION_LIMIT = 50;
|
|
||||||
const DEFAULT_MESSAGE_LIMIT = 50;
|
|
||||||
const MAX_MESSAGE_LIMIT = 200;
|
|
||||||
const STREAM_POLL_INTERVAL_MS = 1000;
|
|
||||||
const INTERNAL_OPERATOR_ID = "internal-provider";
|
|
||||||
|
|
||||||
@Injectable()
|
|
||||||
export class InternalAgentProvider implements IAgentProvider {
|
|
||||||
readonly providerId = "internal";
|
|
||||||
readonly providerType = "internal";
|
|
||||||
readonly displayName = "Internal Orchestrator";
|
|
||||||
|
|
||||||
constructor(
|
|
||||||
private readonly messagesService: AgentMessagesService,
|
|
||||||
private readonly controlService: AgentControlService,
|
|
||||||
private readonly treeService: AgentTreeService
|
|
||||||
) {}
|
|
||||||
|
|
||||||
async listSessions(cursor?: string, limit = DEFAULT_SESSION_LIMIT): Promise<AgentSessionList> {
|
|
||||||
const {
|
|
||||||
sessions,
|
|
||||||
total,
|
|
||||||
cursor: nextCursor,
|
|
||||||
} = await this.treeService.listSessions(cursor, limit);
|
|
||||||
|
|
||||||
return {
|
|
||||||
sessions: sessions.map((session) => this.toAgentSession(session)),
|
|
||||||
total,
|
|
||||||
...(nextCursor !== undefined ? { cursor: nextCursor } : {}),
|
|
||||||
};
|
|
||||||
}
|
|
||||||
|
|
||||||
async getSession(sessionId: string): Promise<AgentSession | null> {
|
|
||||||
const session = await this.treeService.getSession(sessionId);
|
|
||||||
return session ? this.toAgentSession(session) : null;
|
|
||||||
}
|
|
||||||
|
|
||||||
async getMessages(
|
|
||||||
sessionId: string,
|
|
||||||
limit = DEFAULT_MESSAGE_LIMIT,
|
|
||||||
before?: string
|
|
||||||
): Promise<AgentMessage[]> {
|
|
||||||
const safeLimit = this.normalizeMessageLimit(limit);
|
|
||||||
const skip = this.parseSkip(before);
|
|
||||||
|
|
||||||
const result = await this.messagesService.getMessages(sessionId, safeLimit, skip);
|
|
||||||
return result.messages.map((message) => this.toAgentMessage(message));
|
|
||||||
}
|
|
||||||
|
|
||||||
async injectMessage(sessionId: string, content: string): Promise<InjectResult> {
|
|
||||||
await this.controlService.injectMessage(sessionId, INTERNAL_OPERATOR_ID, content);
|
|
||||||
|
|
||||||
return {
|
|
||||||
accepted: true,
|
|
||||||
};
|
|
||||||
}
|
|
||||||
|
|
||||||
async pauseSession(sessionId: string): Promise<void> {
|
|
||||||
await this.controlService.pauseAgent(sessionId, INTERNAL_OPERATOR_ID);
|
|
||||||
}
|
|
||||||
|
|
||||||
async resumeSession(sessionId: string): Promise<void> {
|
|
||||||
await this.controlService.resumeAgent(sessionId, INTERNAL_OPERATOR_ID);
|
|
||||||
}
|
|
||||||
|
|
||||||
async killSession(sessionId: string, force = true): Promise<void> {
|
|
||||||
await this.controlService.killAgent(sessionId, INTERNAL_OPERATOR_ID, force);
|
|
||||||
}
|
|
||||||
|
|
||||||
async *streamMessages(sessionId: string): AsyncIterable<AgentMessage> {
|
|
||||||
const replayMessages = await this.messagesService.getReplayMessages(
|
|
||||||
sessionId,
|
|
||||||
DEFAULT_MESSAGE_LIMIT
|
|
||||||
);
|
|
||||||
|
|
||||||
let lastSeenTimestamp = new Date();
|
|
||||||
let lastSeenMessageId: string | null = null;
|
|
||||||
|
|
||||||
for (const message of replayMessages) {
|
|
||||||
yield this.toAgentMessage(message);
|
|
||||||
lastSeenTimestamp = message.timestamp;
|
|
||||||
lastSeenMessageId = message.id;
|
|
||||||
}
|
|
||||||
|
|
||||||
for (;;) {
|
|
||||||
const newMessages = await this.messagesService.getMessagesAfter(
|
|
||||||
sessionId,
|
|
||||||
lastSeenTimestamp,
|
|
||||||
lastSeenMessageId
|
|
||||||
);
|
|
||||||
|
|
||||||
for (const message of newMessages) {
|
|
||||||
yield this.toAgentMessage(message);
|
|
||||||
lastSeenTimestamp = message.timestamp;
|
|
||||||
lastSeenMessageId = message.id;
|
|
||||||
}
|
|
||||||
|
|
||||||
await this.delay(STREAM_POLL_INTERVAL_MS);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
isAvailable(): Promise<boolean> {
|
|
||||||
return Promise.resolve(true);
|
|
||||||
}
|
|
||||||
|
|
||||||
private toAgentSession(session: AgentSessionTree): AgentSession {
|
|
||||||
const metadata = this.toMetadata(session.metadata);
|
|
||||||
|
|
||||||
return {
|
|
||||||
id: session.sessionId,
|
|
||||||
providerId: this.providerId,
|
|
||||||
providerType: this.providerType,
|
|
||||||
...(session.taskId !== null ? { label: session.taskId } : {}),
|
|
||||||
status: this.toSessionStatus(session.status),
|
|
||||||
...(session.parentSessionId !== null ? { parentSessionId: session.parentSessionId } : {}),
|
|
||||||
createdAt: session.spawnedAt,
|
|
||||||
updatedAt: session.completedAt ?? session.spawnedAt,
|
|
||||||
...(metadata !== undefined ? { metadata } : {}),
|
|
||||||
};
|
|
||||||
}
|
|
||||||
|
|
||||||
private toAgentMessage(message: AgentConversationMessage): AgentMessage {
|
|
||||||
const metadata = this.toMetadata(message.metadata);
|
|
||||||
|
|
||||||
return {
|
|
||||||
id: message.id,
|
|
||||||
sessionId: message.sessionId,
|
|
||||||
role: this.toMessageRole(message.role),
|
|
||||||
content: message.content,
|
|
||||||
timestamp: message.timestamp,
|
|
||||||
...(metadata !== undefined ? { metadata } : {}),
|
|
||||||
};
|
|
||||||
}
|
|
||||||
|
|
||||||
private toSessionStatus(status: string): AgentSessionStatus {
|
|
||||||
switch (status) {
|
|
||||||
case "running":
|
|
||||||
return "active";
|
|
||||||
case "paused":
|
|
||||||
return "paused";
|
|
||||||
case "completed":
|
|
||||||
return "completed";
|
|
||||||
case "failed":
|
|
||||||
case "killed":
|
|
||||||
return "failed";
|
|
||||||
case "spawning":
|
|
||||||
default:
|
|
||||||
return "idle";
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
private toMessageRole(role: string): AgentMessageRole {
|
|
||||||
switch (role) {
|
|
||||||
case "agent":
|
|
||||||
case "assistant":
|
|
||||||
return "assistant";
|
|
||||||
case "system":
|
|
||||||
return "system";
|
|
||||||
case "tool":
|
|
||||||
return "tool";
|
|
||||||
case "operator":
|
|
||||||
case "user":
|
|
||||||
default:
|
|
||||||
return "user";
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
private normalizeMessageLimit(limit: number): number {
|
|
||||||
const normalized = Number.isFinite(limit) ? Math.trunc(limit) : DEFAULT_MESSAGE_LIMIT;
|
|
||||||
if (normalized < 1) {
|
|
||||||
return 1;
|
|
||||||
}
|
|
||||||
|
|
||||||
return Math.min(normalized, MAX_MESSAGE_LIMIT);
|
|
||||||
}
|
|
||||||
|
|
||||||
private parseSkip(before?: string): number {
|
|
||||||
if (!before) {
|
|
||||||
return 0;
|
|
||||||
}
|
|
||||||
|
|
||||||
const parsed = Number.parseInt(before, 10);
|
|
||||||
if (Number.isNaN(parsed) || parsed < 0) {
|
|
||||||
return 0;
|
|
||||||
}
|
|
||||||
|
|
||||||
return parsed;
|
|
||||||
}
|
|
||||||
|
|
||||||
private toMetadata(value: unknown): Record<string, unknown> | undefined {
|
|
||||||
if (value !== null && typeof value === "object" && !Array.isArray(value)) {
|
|
||||||
return value as Record<string, unknown>;
|
|
||||||
}
|
|
||||||
|
|
||||||
return undefined;
|
|
||||||
}
|
|
||||||
|
|
||||||
private async delay(ms: number): Promise<void> {
|
|
||||||
await new Promise((resolve) => {
|
|
||||||
setTimeout(resolve, ms);
|
|
||||||
});
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -1,15 +0,0 @@
|
|||||||
import { Type } from "class-transformer";
|
|
||||||
import { IsInt, IsOptional, IsString, Max, Min } from "class-validator";
|
|
||||||
|
|
||||||
export class GetMissionControlMessagesQueryDto {
|
|
||||||
@IsOptional()
|
|
||||||
@Type(() => Number)
|
|
||||||
@IsInt()
|
|
||||||
@Min(1)
|
|
||||||
@Max(200)
|
|
||||||
limit?: number;
|
|
||||||
|
|
||||||
@IsOptional()
|
|
||||||
@IsString()
|
|
||||||
before?: string;
|
|
||||||
}
|
|
||||||
@@ -1,7 +0,0 @@
|
|||||||
import { IsBoolean, IsOptional } from "class-validator";
|
|
||||||
|
|
||||||
export class KillSessionDto {
|
|
||||||
@IsOptional()
|
|
||||||
@IsBoolean()
|
|
||||||
force?: boolean;
|
|
||||||
}
|
|
||||||
@@ -1,183 +0,0 @@
|
|||||||
import {
|
|
||||||
Body,
|
|
||||||
Controller,
|
|
||||||
Get,
|
|
||||||
Header,
|
|
||||||
HttpCode,
|
|
||||||
MessageEvent,
|
|
||||||
Param,
|
|
||||||
Post,
|
|
||||||
Query,
|
|
||||||
Request,
|
|
||||||
Sse,
|
|
||||||
UseGuards,
|
|
||||||
UsePipes,
|
|
||||||
ValidationPipe,
|
|
||||||
} from "@nestjs/common";
|
|
||||||
import type { AgentMessage, AgentSession, InjectResult } from "@mosaic/shared";
|
|
||||||
import { Observable } from "rxjs";
|
|
||||||
import { AuthGuard } from "../../auth/guards/auth.guard";
|
|
||||||
import { InjectAgentDto } from "../agents/dto/inject-agent.dto";
|
|
||||||
import { GetMissionControlMessagesQueryDto } from "./dto/get-mission-control-messages-query.dto";
|
|
||||||
import { KillSessionDto } from "./dto/kill-session.dto";
|
|
||||||
import { MissionControlService } from "./mission-control.service";
|
|
||||||
|
|
||||||
const DEFAULT_OPERATOR_ID = "mission-control";
|
|
||||||
|
|
||||||
interface MissionControlRequest {
|
|
||||||
user?: {
|
|
||||||
id?: string;
|
|
||||||
};
|
|
||||||
}
|
|
||||||
|
|
||||||
@Controller("api/mission-control")
|
|
||||||
@UseGuards(AuthGuard)
|
|
||||||
export class MissionControlController {
|
|
||||||
constructor(private readonly missionControlService: MissionControlService) {}
|
|
||||||
|
|
||||||
@Get("sessions")
|
|
||||||
async listSessions(): Promise<{ sessions: AgentSession[] }> {
|
|
||||||
const sessions = await this.missionControlService.listSessions();
|
|
||||||
return { sessions };
|
|
||||||
}
|
|
||||||
|
|
||||||
@Get("sessions/:sessionId")
|
|
||||||
getSession(@Param("sessionId") sessionId: string): Promise<AgentSession> {
|
|
||||||
return this.missionControlService.getSession(sessionId);
|
|
||||||
}
|
|
||||||
|
|
||||||
@Get("sessions/:sessionId/messages")
|
|
||||||
@UsePipes(new ValidationPipe({ transform: true, whitelist: true }))
|
|
||||||
async getMessages(
|
|
||||||
@Param("sessionId") sessionId: string,
|
|
||||||
@Query() query: GetMissionControlMessagesQueryDto
|
|
||||||
): Promise<{ messages: AgentMessage[] }> {
|
|
||||||
const messages = await this.missionControlService.getMessages(
|
|
||||||
sessionId,
|
|
||||||
query.limit,
|
|
||||||
query.before
|
|
||||||
);
|
|
||||||
|
|
||||||
return { messages };
|
|
||||||
}
|
|
||||||
|
|
||||||
@Post("sessions/:sessionId/inject")
|
|
||||||
@HttpCode(200)
|
|
||||||
@UsePipes(new ValidationPipe({ transform: true, whitelist: true }))
|
|
||||||
injectMessage(
|
|
||||||
@Param("sessionId") sessionId: string,
|
|
||||||
@Body() dto: InjectAgentDto,
|
|
||||||
@Request() req: MissionControlRequest
|
|
||||||
): Promise<InjectResult> {
|
|
||||||
return this.missionControlService.injectMessage(
|
|
||||||
sessionId,
|
|
||||||
dto.message,
|
|
||||||
this.resolveOperatorId(req)
|
|
||||||
);
|
|
||||||
}
|
|
||||||
|
|
||||||
@Post("sessions/:sessionId/pause")
|
|
||||||
@HttpCode(200)
|
|
||||||
async pauseSession(
|
|
||||||
@Param("sessionId") sessionId: string,
|
|
||||||
@Request() req: MissionControlRequest
|
|
||||||
): Promise<{ message: string }> {
|
|
||||||
await this.missionControlService.pauseSession(sessionId, this.resolveOperatorId(req));
|
|
||||||
|
|
||||||
return { message: `Session ${sessionId} paused` };
|
|
||||||
}
|
|
||||||
|
|
||||||
@Post("sessions/:sessionId/resume")
|
|
||||||
@HttpCode(200)
|
|
||||||
async resumeSession(
|
|
||||||
@Param("sessionId") sessionId: string,
|
|
||||||
@Request() req: MissionControlRequest
|
|
||||||
): Promise<{ message: string }> {
|
|
||||||
await this.missionControlService.resumeSession(sessionId, this.resolveOperatorId(req));
|
|
||||||
|
|
||||||
return { message: `Session ${sessionId} resumed` };
|
|
||||||
}
|
|
||||||
|
|
||||||
@Post("sessions/:sessionId/kill")
|
|
||||||
@HttpCode(200)
|
|
||||||
@UsePipes(new ValidationPipe({ transform: true, whitelist: true }))
|
|
||||||
async killSession(
|
|
||||||
@Param("sessionId") sessionId: string,
|
|
||||||
@Body() dto: KillSessionDto,
|
|
||||||
@Request() req: MissionControlRequest
|
|
||||||
): Promise<{ message: string }> {
|
|
||||||
await this.missionControlService.killSession(
|
|
||||||
sessionId,
|
|
||||||
dto.force ?? true,
|
|
||||||
this.resolveOperatorId(req)
|
|
||||||
);
|
|
||||||
|
|
||||||
return { message: `Session ${sessionId} killed` };
|
|
||||||
}
|
|
||||||
|
|
||||||
@Sse("sessions/:sessionId/stream")
|
|
||||||
@Header("Content-Type", "text/event-stream")
|
|
||||||
@Header("Cache-Control", "no-cache")
|
|
||||||
streamSessionMessages(@Param("sessionId") sessionId: string): Observable<MessageEvent> {
|
|
||||||
return new Observable<MessageEvent>((subscriber) => {
|
|
||||||
let isClosed = false;
|
|
||||||
let iterator: AsyncIterator<AgentMessage> | null = null;
|
|
||||||
|
|
||||||
void this.missionControlService
|
|
||||||
.streamMessages(sessionId)
|
|
||||||
.then(async (stream) => {
|
|
||||||
iterator = stream[Symbol.asyncIterator]();
|
|
||||||
|
|
||||||
for (;;) {
|
|
||||||
if (isClosed) {
|
|
||||||
break;
|
|
||||||
}
|
|
||||||
|
|
||||||
const next = (await iterator.next()) as { done: boolean; value: AgentMessage };
|
|
||||||
if (next.done) {
|
|
||||||
break;
|
|
||||||
}
|
|
||||||
|
|
||||||
subscriber.next({
|
|
||||||
data: this.toStreamPayload(next.value),
|
|
||||||
});
|
|
||||||
}
|
|
||||||
|
|
||||||
subscriber.complete();
|
|
||||||
})
|
|
||||||
.catch((error: unknown) => {
|
|
||||||
subscriber.error(error);
|
|
||||||
});
|
|
||||||
|
|
||||||
return () => {
|
|
||||||
isClosed = true;
|
|
||||||
void iterator?.return?.();
|
|
||||||
};
|
|
||||||
});
|
|
||||||
}
|
|
||||||
|
|
||||||
private resolveOperatorId(req: MissionControlRequest): string {
|
|
||||||
const operatorId = req.user?.id;
|
|
||||||
return typeof operatorId === "string" && operatorId.length > 0
|
|
||||||
? operatorId
|
|
||||||
: DEFAULT_OPERATOR_ID;
|
|
||||||
}
|
|
||||||
|
|
||||||
private toStreamPayload(message: AgentMessage): {
|
|
||||||
id: string;
|
|
||||||
sessionId: string;
|
|
||||||
role: string;
|
|
||||||
content: string;
|
|
||||||
timestamp: string;
|
|
||||||
metadata?: Record<string, unknown>;
|
|
||||||
} {
|
|
||||||
return {
|
|
||||||
id: message.id,
|
|
||||||
sessionId: message.sessionId,
|
|
||||||
role: message.role,
|
|
||||||
content: message.content,
|
|
||||||
timestamp: message.timestamp.toISOString(),
|
|
||||||
...(message.metadata !== undefined ? { metadata: message.metadata } : {}),
|
|
||||||
};
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -1,13 +0,0 @@
|
|||||||
import { Module } from "@nestjs/common";
|
|
||||||
import { AgentsModule } from "../agents/agents.module";
|
|
||||||
import { AuthModule } from "../../auth/auth.module";
|
|
||||||
import { PrismaModule } from "../../prisma/prisma.module";
|
|
||||||
import { MissionControlController } from "./mission-control.controller";
|
|
||||||
import { MissionControlService } from "./mission-control.service";
|
|
||||||
|
|
||||||
@Module({
|
|
||||||
imports: [AgentsModule, AuthModule, PrismaModule],
|
|
||||||
controllers: [MissionControlController],
|
|
||||||
providers: [MissionControlService],
|
|
||||||
})
|
|
||||||
export class MissionControlModule {}
|
|
||||||
@@ -1,213 +0,0 @@
|
|||||||
import { NotFoundException } from "@nestjs/common";
|
|
||||||
import { beforeEach, describe, expect, it, vi } from "vitest";
|
|
||||||
import type { AgentMessage, AgentSession, IAgentProvider, InjectResult } from "@mosaic/shared";
|
|
||||||
import type { PrismaService } from "../../prisma/prisma.service";
|
|
||||||
import { AgentProviderRegistry } from "../agents/agent-provider.registry";
|
|
||||||
import { MissionControlService } from "./mission-control.service";
|
|
||||||
|
|
||||||
type MockProvider = IAgentProvider & {
|
|
||||||
listSessions: ReturnType<typeof vi.fn>;
|
|
||||||
getSession: ReturnType<typeof vi.fn>;
|
|
||||||
getMessages: ReturnType<typeof vi.fn>;
|
|
||||||
injectMessage: ReturnType<typeof vi.fn>;
|
|
||||||
pauseSession: ReturnType<typeof vi.fn>;
|
|
||||||
resumeSession: ReturnType<typeof vi.fn>;
|
|
||||||
killSession: ReturnType<typeof vi.fn>;
|
|
||||||
streamMessages: ReturnType<typeof vi.fn>;
|
|
||||||
};
|
|
||||||
|
|
||||||
const emptyMessageStream = async function* (): AsyncIterable<AgentMessage> {
|
|
||||||
return;
|
|
||||||
};
|
|
||||||
|
|
||||||
const createProvider = (providerId = "internal"): MockProvider => ({
|
|
||||||
providerId,
|
|
||||||
providerType: providerId,
|
|
||||||
displayName: providerId,
|
|
||||||
listSessions: vi.fn().mockResolvedValue({ sessions: [], total: 0 }),
|
|
||||||
getSession: vi.fn().mockResolvedValue(null),
|
|
||||||
getMessages: vi.fn().mockResolvedValue([]),
|
|
||||||
injectMessage: vi.fn().mockResolvedValue({ accepted: true } as InjectResult),
|
|
||||||
pauseSession: vi.fn().mockResolvedValue(undefined),
|
|
||||||
resumeSession: vi.fn().mockResolvedValue(undefined),
|
|
||||||
killSession: vi.fn().mockResolvedValue(undefined),
|
|
||||||
streamMessages: vi.fn().mockReturnValue(emptyMessageStream()),
|
|
||||||
isAvailable: vi.fn().mockResolvedValue(true),
|
|
||||||
});
|
|
||||||
|
|
||||||
describe("MissionControlService", () => {
|
|
||||||
let service: MissionControlService;
|
|
||||||
let registry: {
|
|
||||||
listAllSessions: ReturnType<typeof vi.fn>;
|
|
||||||
getProviderForSession: ReturnType<typeof vi.fn>;
|
|
||||||
};
|
|
||||||
let prisma: {
|
|
||||||
operatorAuditLog: {
|
|
||||||
create: ReturnType<typeof vi.fn>;
|
|
||||||
};
|
|
||||||
};
|
|
||||||
|
|
||||||
const session: AgentSession = {
|
|
||||||
id: "session-1",
|
|
||||||
providerId: "internal",
|
|
||||||
providerType: "internal",
|
|
||||||
status: "active",
|
|
||||||
createdAt: new Date("2026-03-07T14:00:00.000Z"),
|
|
||||||
updatedAt: new Date("2026-03-07T14:01:00.000Z"),
|
|
||||||
};
|
|
||||||
|
|
||||||
beforeEach(() => {
|
|
||||||
registry = {
|
|
||||||
listAllSessions: vi.fn().mockResolvedValue([session]),
|
|
||||||
getProviderForSession: vi.fn().mockResolvedValue(null),
|
|
||||||
};
|
|
||||||
|
|
||||||
prisma = {
|
|
||||||
operatorAuditLog: {
|
|
||||||
create: vi.fn().mockResolvedValue(undefined),
|
|
||||||
},
|
|
||||||
};
|
|
||||||
|
|
||||||
service = new MissionControlService(
|
|
||||||
registry as unknown as AgentProviderRegistry,
|
|
||||||
prisma as unknown as PrismaService
|
|
||||||
);
|
|
||||||
});
|
|
||||||
|
|
||||||
it("lists sessions from the registry", async () => {
|
|
||||||
await expect(service.listSessions()).resolves.toEqual([session]);
|
|
||||||
expect(registry.listAllSessions).toHaveBeenCalledTimes(1);
|
|
||||||
});
|
|
||||||
|
|
||||||
it("returns a session when it is found", async () => {
|
|
||||||
const provider = createProvider("internal");
|
|
||||||
registry.getProviderForSession.mockResolvedValue({ provider, session });
|
|
||||||
|
|
||||||
await expect(service.getSession(session.id)).resolves.toEqual(session);
|
|
||||||
});
|
|
||||||
|
|
||||||
it("throws NotFoundException when session lookup fails", async () => {
|
|
||||||
await expect(service.getSession("missing-session")).rejects.toBeInstanceOf(NotFoundException);
|
|
||||||
});
|
|
||||||
|
|
||||||
it("gets messages from the resolved provider", async () => {
|
|
||||||
const provider = createProvider("openclaw");
|
|
||||||
const messages: AgentMessage[] = [
|
|
||||||
{
|
|
||||||
id: "message-1",
|
|
||||||
sessionId: session.id,
|
|
||||||
role: "assistant",
|
|
||||||
content: "hello",
|
|
||||||
timestamp: new Date("2026-03-07T14:01:00.000Z"),
|
|
||||||
},
|
|
||||||
];
|
|
||||||
|
|
||||||
provider.getMessages.mockResolvedValue(messages);
|
|
||||||
registry.getProviderForSession.mockResolvedValue({ provider, session });
|
|
||||||
|
|
||||||
await expect(service.getMessages(session.id, 25, "10")).resolves.toEqual(messages);
|
|
||||||
expect(provider.getMessages).toHaveBeenCalledWith(session.id, 25, "10");
|
|
||||||
});
|
|
||||||
|
|
||||||
it("injects a message and writes an audit log", async () => {
|
|
||||||
const provider = createProvider("internal");
|
|
||||||
const injectResult: InjectResult = { accepted: true, messageId: "msg-1" };
|
|
||||||
provider.injectMessage.mockResolvedValue(injectResult);
|
|
||||||
registry.getProviderForSession.mockResolvedValue({ provider, session });
|
|
||||||
|
|
||||||
await expect(service.injectMessage(session.id, "ship it", "operator-1")).resolves.toEqual(
|
|
||||||
injectResult
|
|
||||||
);
|
|
||||||
|
|
||||||
expect(provider.injectMessage).toHaveBeenCalledWith(session.id, "ship it");
|
|
||||||
expect(prisma.operatorAuditLog.create).toHaveBeenCalledWith({
|
|
||||||
data: {
|
|
||||||
sessionId: session.id,
|
|
||||||
userId: "operator-1",
|
|
||||||
provider: "internal",
|
|
||||||
action: "inject",
|
|
||||||
content: "ship it",
|
|
||||||
metadata: {
|
|
||||||
payload: { message: "ship it" },
|
|
||||||
},
|
|
||||||
},
|
|
||||||
});
|
|
||||||
});
|
|
||||||
|
|
||||||
it("pauses and resumes using default operator id", async () => {
|
|
||||||
const provider = createProvider("openclaw");
|
|
||||||
registry.getProviderForSession.mockResolvedValue({ provider, session });
|
|
||||||
|
|
||||||
await service.pauseSession(session.id);
|
|
||||||
await service.resumeSession(session.id);
|
|
||||||
|
|
||||||
expect(provider.pauseSession).toHaveBeenCalledWith(session.id);
|
|
||||||
expect(provider.resumeSession).toHaveBeenCalledWith(session.id);
|
|
||||||
expect(prisma.operatorAuditLog.create).toHaveBeenNthCalledWith(1, {
|
|
||||||
data: {
|
|
||||||
sessionId: session.id,
|
|
||||||
userId: "mission-control",
|
|
||||||
provider: "openclaw",
|
|
||||||
action: "pause",
|
|
||||||
metadata: {
|
|
||||||
payload: {},
|
|
||||||
},
|
|
||||||
},
|
|
||||||
});
|
|
||||||
expect(prisma.operatorAuditLog.create).toHaveBeenNthCalledWith(2, {
|
|
||||||
data: {
|
|
||||||
sessionId: session.id,
|
|
||||||
userId: "mission-control",
|
|
||||||
provider: "openclaw",
|
|
||||||
action: "resume",
|
|
||||||
metadata: {
|
|
||||||
payload: {},
|
|
||||||
},
|
|
||||||
},
|
|
||||||
});
|
|
||||||
});
|
|
||||||
|
|
||||||
it("kills with provided force value and writes audit log", async () => {
|
|
||||||
const provider = createProvider("openclaw");
|
|
||||||
registry.getProviderForSession.mockResolvedValue({ provider, session });
|
|
||||||
|
|
||||||
await service.killSession(session.id, false, "operator-2");
|
|
||||||
|
|
||||||
expect(provider.killSession).toHaveBeenCalledWith(session.id, false);
|
|
||||||
expect(prisma.operatorAuditLog.create).toHaveBeenCalledWith({
|
|
||||||
data: {
|
|
||||||
sessionId: session.id,
|
|
||||||
userId: "operator-2",
|
|
||||||
provider: "openclaw",
|
|
||||||
action: "kill",
|
|
||||||
metadata: {
|
|
||||||
payload: { force: false },
|
|
||||||
},
|
|
||||||
},
|
|
||||||
});
|
|
||||||
});
|
|
||||||
|
|
||||||
it("resolves provider message stream", async () => {
|
|
||||||
const provider = createProvider("internal");
|
|
||||||
const messageStream = (async function* (): AsyncIterable<AgentMessage> {
|
|
||||||
yield {
|
|
||||||
id: "message-1",
|
|
||||||
sessionId: session.id,
|
|
||||||
role: "assistant",
|
|
||||||
content: "stream",
|
|
||||||
timestamp: new Date("2026-03-07T14:03:00.000Z"),
|
|
||||||
};
|
|
||||||
})();
|
|
||||||
|
|
||||||
provider.streamMessages.mockReturnValue(messageStream);
|
|
||||||
registry.getProviderForSession.mockResolvedValue({ provider, session });
|
|
||||||
|
|
||||||
await expect(service.streamMessages(session.id)).resolves.toBe(messageStream);
|
|
||||||
expect(provider.streamMessages).toHaveBeenCalledWith(session.id);
|
|
||||||
});
|
|
||||||
|
|
||||||
it("does not write audit log when session cannot be resolved", async () => {
|
|
||||||
await expect(service.pauseSession("missing-session")).rejects.toBeInstanceOf(NotFoundException);
|
|
||||||
expect(prisma.operatorAuditLog.create).not.toHaveBeenCalled();
|
|
||||||
});
|
|
||||||
});
|
|
||||||
@@ -1,139 +0,0 @@
|
|||||||
import { Injectable, NotFoundException } from "@nestjs/common";
|
|
||||||
import type { AgentMessage, AgentSession, IAgentProvider, InjectResult } from "@mosaic/shared";
|
|
||||||
import type { Prisma } from "@prisma/client";
|
|
||||||
import { PrismaService } from "../../prisma/prisma.service";
|
|
||||||
import { AgentProviderRegistry } from "../agents/agent-provider.registry";
|
|
||||||
|
|
||||||
type MissionControlAction = "inject" | "pause" | "resume" | "kill";
|
|
||||||
|
|
||||||
const DEFAULT_OPERATOR_ID = "mission-control";
|
|
||||||
|
|
||||||
@Injectable()
|
|
||||||
export class MissionControlService {
|
|
||||||
constructor(
|
|
||||||
private readonly registry: AgentProviderRegistry,
|
|
||||||
private readonly prisma: PrismaService
|
|
||||||
) {}
|
|
||||||
|
|
||||||
listSessions(): Promise<AgentSession[]> {
|
|
||||||
return this.registry.listAllSessions();
|
|
||||||
}
|
|
||||||
|
|
||||||
async getSession(sessionId: string): Promise<AgentSession> {
|
|
||||||
const resolved = await this.registry.getProviderForSession(sessionId);
|
|
||||||
if (!resolved) {
|
|
||||||
throw new NotFoundException(`Session ${sessionId} not found`);
|
|
||||||
}
|
|
||||||
|
|
||||||
return resolved.session;
|
|
||||||
}
|
|
||||||
|
|
||||||
async getMessages(sessionId: string, limit?: number, before?: string): Promise<AgentMessage[]> {
|
|
||||||
const { provider } = await this.getProviderForSessionOrThrow(sessionId);
|
|
||||||
return provider.getMessages(sessionId, limit, before);
|
|
||||||
}
|
|
||||||
|
|
||||||
async injectMessage(
|
|
||||||
sessionId: string,
|
|
||||||
message: string,
|
|
||||||
operatorId = DEFAULT_OPERATOR_ID
|
|
||||||
): Promise<InjectResult> {
|
|
||||||
const { provider } = await this.getProviderForSessionOrThrow(sessionId);
|
|
||||||
const result = await provider.injectMessage(sessionId, message);
|
|
||||||
|
|
||||||
await this.writeOperatorAuditLog({
|
|
||||||
sessionId,
|
|
||||||
providerId: provider.providerId,
|
|
||||||
operatorId,
|
|
||||||
action: "inject",
|
|
||||||
content: message,
|
|
||||||
payload: { message },
|
|
||||||
});
|
|
||||||
|
|
||||||
return result;
|
|
||||||
}
|
|
||||||
|
|
||||||
async pauseSession(sessionId: string, operatorId = DEFAULT_OPERATOR_ID): Promise<void> {
|
|
||||||
const { provider } = await this.getProviderForSessionOrThrow(sessionId);
|
|
||||||
await provider.pauseSession(sessionId);
|
|
||||||
|
|
||||||
await this.writeOperatorAuditLog({
|
|
||||||
sessionId,
|
|
||||||
providerId: provider.providerId,
|
|
||||||
operatorId,
|
|
||||||
action: "pause",
|
|
||||||
payload: {},
|
|
||||||
});
|
|
||||||
}
|
|
||||||
|
|
||||||
async resumeSession(sessionId: string, operatorId = DEFAULT_OPERATOR_ID): Promise<void> {
|
|
||||||
const { provider } = await this.getProviderForSessionOrThrow(sessionId);
|
|
||||||
await provider.resumeSession(sessionId);
|
|
||||||
|
|
||||||
await this.writeOperatorAuditLog({
|
|
||||||
sessionId,
|
|
||||||
providerId: provider.providerId,
|
|
||||||
operatorId,
|
|
||||||
action: "resume",
|
|
||||||
payload: {},
|
|
||||||
});
|
|
||||||
}
|
|
||||||
|
|
||||||
async killSession(
|
|
||||||
sessionId: string,
|
|
||||||
force = true,
|
|
||||||
operatorId = DEFAULT_OPERATOR_ID
|
|
||||||
): Promise<void> {
|
|
||||||
const { provider } = await this.getProviderForSessionOrThrow(sessionId);
|
|
||||||
await provider.killSession(sessionId, force);
|
|
||||||
|
|
||||||
await this.writeOperatorAuditLog({
|
|
||||||
sessionId,
|
|
||||||
providerId: provider.providerId,
|
|
||||||
operatorId,
|
|
||||||
action: "kill",
|
|
||||||
payload: { force },
|
|
||||||
});
|
|
||||||
}
|
|
||||||
|
|
||||||
async streamMessages(sessionId: string): Promise<AsyncIterable<AgentMessage>> {
|
|
||||||
const { provider } = await this.getProviderForSessionOrThrow(sessionId);
|
|
||||||
return provider.streamMessages(sessionId);
|
|
||||||
}
|
|
||||||
|
|
||||||
private async getProviderForSessionOrThrow(
|
|
||||||
sessionId: string
|
|
||||||
): Promise<{ provider: IAgentProvider; session: AgentSession }> {
|
|
||||||
const resolved = await this.registry.getProviderForSession(sessionId);
|
|
||||||
|
|
||||||
if (!resolved) {
|
|
||||||
throw new NotFoundException(`Session ${sessionId} not found`);
|
|
||||||
}
|
|
||||||
|
|
||||||
return resolved;
|
|
||||||
}
|
|
||||||
|
|
||||||
private toJsonValue(value: Record<string, unknown>): Prisma.InputJsonValue {
|
|
||||||
return value as Prisma.InputJsonValue;
|
|
||||||
}
|
|
||||||
|
|
||||||
private async writeOperatorAuditLog(params: {
|
|
||||||
sessionId: string;
|
|
||||||
providerId: string;
|
|
||||||
operatorId: string;
|
|
||||||
action: MissionControlAction;
|
|
||||||
content?: string;
|
|
||||||
payload: Record<string, unknown>;
|
|
||||||
}): Promise<void> {
|
|
||||||
await this.prisma.operatorAuditLog.create({
|
|
||||||
data: {
|
|
||||||
sessionId: params.sessionId,
|
|
||||||
userId: params.operatorId,
|
|
||||||
provider: params.providerId,
|
|
||||||
action: params.action,
|
|
||||||
...(params.content !== undefined ? { content: params.content } : {}),
|
|
||||||
metadata: this.toJsonValue({ payload: params.payload }),
|
|
||||||
},
|
|
||||||
});
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -4,9 +4,7 @@ import { BullModule } from "@nestjs/bullmq";
|
|||||||
import { ThrottlerModule } from "@nestjs/throttler";
|
import { ThrottlerModule } from "@nestjs/throttler";
|
||||||
import { HealthModule } from "./api/health/health.module";
|
import { HealthModule } from "./api/health/health.module";
|
||||||
import { AgentsModule } from "./api/agents/agents.module";
|
import { AgentsModule } from "./api/agents/agents.module";
|
||||||
import { MissionControlModule } from "./api/mission-control/mission-control.module";
|
|
||||||
import { QueueApiModule } from "./api/queue/queue-api.module";
|
import { QueueApiModule } from "./api/queue/queue-api.module";
|
||||||
import { AgentProvidersModule } from "./api/agent-providers/agent-providers.module";
|
|
||||||
import { CoordinatorModule } from "./coordinator/coordinator.module";
|
import { CoordinatorModule } from "./coordinator/coordinator.module";
|
||||||
import { BudgetModule } from "./budget/budget.module";
|
import { BudgetModule } from "./budget/budget.module";
|
||||||
import { CIModule } from "./ci";
|
import { CIModule } from "./ci";
|
||||||
@@ -53,8 +51,6 @@ import { orchestratorConfig } from "./config/orchestrator.config";
|
|||||||
]),
|
]),
|
||||||
HealthModule,
|
HealthModule,
|
||||||
AgentsModule,
|
AgentsModule,
|
||||||
AgentProvidersModule,
|
|
||||||
MissionControlModule,
|
|
||||||
QueueApiModule,
|
QueueApiModule,
|
||||||
CoordinatorModule,
|
CoordinatorModule,
|
||||||
BudgetModule,
|
BudgetModule,
|
||||||
|
|||||||
@@ -1,9 +0,0 @@
|
|||||||
import { Module } from "@nestjs/common";
|
|
||||||
import { OrchestratorApiKeyGuard } from "../common/guards/api-key.guard";
|
|
||||||
import { AuthGuard } from "./guards/auth.guard";
|
|
||||||
|
|
||||||
@Module({
|
|
||||||
providers: [OrchestratorApiKeyGuard, AuthGuard],
|
|
||||||
exports: [AuthGuard],
|
|
||||||
})
|
|
||||||
export class AuthModule {}
|
|
||||||
@@ -1,11 +0,0 @@
|
|||||||
import { CanActivate, ExecutionContext, Injectable } from "@nestjs/common";
|
|
||||||
import { OrchestratorApiKeyGuard } from "../../common/guards/api-key.guard";
|
|
||||||
|
|
||||||
@Injectable()
|
|
||||||
export class AuthGuard implements CanActivate {
|
|
||||||
constructor(private readonly apiKeyGuard: OrchestratorApiKeyGuard) {}
|
|
||||||
|
|
||||||
canActivate(context: ExecutionContext): boolean | Promise<boolean> {
|
|
||||||
return this.apiKeyGuard.canActivate(context);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -115,13 +115,7 @@ export class AgentSpawnerService implements OnModuleDestroy {
|
|||||||
}
|
}
|
||||||
|
|
||||||
void this.agentIngestionService
|
void this.agentIngestionService
|
||||||
.recordAgentSpawned(
|
.recordAgentSpawned(agentId, undefined, undefined, request.taskId, request.agentType)
|
||||||
agentId,
|
|
||||||
request.parentAgentId,
|
|
||||||
undefined,
|
|
||||||
request.taskId,
|
|
||||||
request.agentType
|
|
||||||
)
|
|
||||||
.catch((error: unknown) => {
|
.catch((error: unknown) => {
|
||||||
const errorMessage = error instanceof Error ? error.message : String(error);
|
const errorMessage = error instanceof Error ? error.message : String(error);
|
||||||
this.logger.error(`Failed to record spawned ingestion for ${agentId}: ${errorMessage}`);
|
this.logger.error(`Failed to record spawned ingestion for ${agentId}: ${errorMessage}`);
|
||||||
|
|||||||
@@ -40,8 +40,6 @@ export interface SpawnAgentOptions {
|
|||||||
export interface SpawnAgentRequest {
|
export interface SpawnAgentRequest {
|
||||||
/** Unique task identifier */
|
/** Unique task identifier */
|
||||||
taskId: string;
|
taskId: string;
|
||||||
/** Optional parent session identifier for subagent lineage */
|
|
||||||
parentAgentId?: string;
|
|
||||||
/** Type of agent to spawn */
|
/** Type of agent to spawn */
|
||||||
agentType: AgentType;
|
agentType: AgentType;
|
||||||
/** Context for task execution */
|
/** Context for task execution */
|
||||||
|
|||||||
@@ -124,9 +124,9 @@ Target version: `v0.0.23`
|
|||||||
| MS23-P0-001 | done | p0-foundation | Prisma schema: AgentConversationMessage, AgentSessionTree, AgentProviderConfig, OperatorAuditLog | #693 | api | feat/ms23-p0-schema | — | MS23-P0-002,MS23-P0-003,MS23-P0-004,MS23-P0-005,MS23-P1-001 | codex | 2026-03-06 | 2026-03-06 | 15K | — | taskSource field per mosaic-queue note in PRD |
|
| MS23-P0-001 | done | p0-foundation | Prisma schema: AgentConversationMessage, AgentSessionTree, AgentProviderConfig, OperatorAuditLog | #693 | api | feat/ms23-p0-schema | — | MS23-P0-002,MS23-P0-003,MS23-P0-004,MS23-P0-005,MS23-P1-001 | codex | 2026-03-06 | 2026-03-06 | 15K | — | taskSource field per mosaic-queue note in PRD |
|
||||||
| MS23-P0-002 | done | p0-foundation | Agent message ingestion: wire spawner/lifecycle to write messages to DB | #693 | orchestrator | feat/ms23-p0-ingestion | MS23-P0-001 | MS23-P0-006 | codex | 2026-03-06 | 2026-03-07 | 20K | — | |
|
| MS23-P0-002 | done | p0-foundation | Agent message ingestion: wire spawner/lifecycle to write messages to DB | #693 | orchestrator | feat/ms23-p0-ingestion | MS23-P0-001 | MS23-P0-006 | codex | 2026-03-06 | 2026-03-07 | 20K | — | |
|
||||||
| MS23-P0-003 | done | p0-foundation | Orchestrator API: GET /agents/:id/messages + SSE stream endpoint | #693 | orchestrator | feat/ms23-p0-stream | MS23-P0-001 | MS23-P0-006 | codex | 2026-03-06 | 2026-03-07 | 20K | — | |
|
| MS23-P0-003 | done | p0-foundation | Orchestrator API: GET /agents/:id/messages + SSE stream endpoint | #693 | orchestrator | feat/ms23-p0-stream | MS23-P0-001 | MS23-P0-006 | codex | 2026-03-06 | 2026-03-07 | 20K | — | |
|
||||||
| MS23-P0-004 | done | p0-foundation | Orchestrator API: POST /agents/:id/inject + pause/resume endpoints | #693 | orchestrator | feat/ms23-p0-controls | MS23-P0-001 | MS23-P0-006 | codex | 2026-03-07 | 2026-03-07 | 15K | — | |
|
| MS23-P0-004 | in-progress | p0-foundation | Orchestrator API: POST /agents/:id/inject + pause/resume endpoints | #693 | orchestrator | feat/ms23-p0-controls | MS23-P0-001 | MS23-P0-006 | codex | 2026-03-07 | — | 15K | — | |
|
||||||
| MS23-P0-005 | done | p0-foundation | Subagent tree: parentAgentId on spawn registration + GET /agents/tree | #693 | orchestrator | feat/ms23-p0-tree | MS23-P0-001 | MS23-P0-006 | — | — | — | 15K | — | |
|
| MS23-P0-005 | not-started | p0-foundation | Subagent tree: parentAgentId on spawn registration + GET /agents/tree | #693 | orchestrator | feat/ms23-p0-tree | MS23-P0-001 | MS23-P0-006 | — | — | — | 15K | — | |
|
||||||
| MS23-P0-006 | done | p0-foundation | Unit + integration tests for all P0 orchestrator endpoints | #693 | orchestrator | test/ms23-p0 | MS23-P0-002,MS23-P0-003,MS23-P0-004,MS23-P0-005 | MS23-P1-001 | codex | 2026-03-07 | 2026-03-07 | 20K | — | Phase 0 gate: SSE stream verified via curl |
|
| MS23-P0-006 | not-started | p0-foundation | Unit + integration tests for all P0 orchestrator endpoints | #693 | orchestrator | test/ms23-p0 | MS23-P0-002,MS23-P0-003,MS23-P0-004,MS23-P0-005 | MS23-P1-001 | — | — | — | 20K | — | Phase 0 gate: SSE stream verified via curl |
|
||||||
|
|
||||||
### Phase 1 — Provider Interface (Plugin Architecture)
|
### Phase 1 — Provider Interface (Plugin Architecture)
|
||||||
|
|
||||||
|
|||||||
@@ -1,78 +0,0 @@
|
|||||||
// Agent message roles
|
|
||||||
export type AgentMessageRole = "user" | "assistant" | "system" | "tool";
|
|
||||||
|
|
||||||
// A single message in an agent conversation
|
|
||||||
export interface AgentMessage {
|
|
||||||
id: string;
|
|
||||||
sessionId: string;
|
|
||||||
role: AgentMessageRole;
|
|
||||||
content: string;
|
|
||||||
timestamp: Date;
|
|
||||||
metadata?: Record<string, unknown>;
|
|
||||||
}
|
|
||||||
|
|
||||||
// Session lifecycle status
|
|
||||||
export type AgentSessionStatus = "active" | "paused" | "completed" | "failed" | "idle";
|
|
||||||
|
|
||||||
// An agent session (conversation thread)
|
|
||||||
export interface AgentSession {
|
|
||||||
id: string;
|
|
||||||
providerId: string; // which provider owns this session
|
|
||||||
providerType: string; // "internal" | "openclaw" | etc.
|
|
||||||
label?: string;
|
|
||||||
status: AgentSessionStatus;
|
|
||||||
parentSessionId?: string; // for subagent trees
|
|
||||||
createdAt: Date;
|
|
||||||
updatedAt: Date;
|
|
||||||
metadata?: Record<string, unknown>;
|
|
||||||
}
|
|
||||||
|
|
||||||
// Result of listing sessions
|
|
||||||
export interface AgentSessionList {
|
|
||||||
sessions: AgentSession[];
|
|
||||||
total: number;
|
|
||||||
cursor?: string;
|
|
||||||
}
|
|
||||||
|
|
||||||
// Result of injecting a message
|
|
||||||
export interface InjectResult {
|
|
||||||
accepted: boolean;
|
|
||||||
messageId?: string;
|
|
||||||
}
|
|
||||||
|
|
||||||
// The IAgentProvider interface — every provider (internal, OpenClaw, future) implements this
|
|
||||||
export interface IAgentProvider {
|
|
||||||
readonly providerId: string;
|
|
||||||
readonly providerType: string;
|
|
||||||
readonly displayName: string;
|
|
||||||
|
|
||||||
// Session management
|
|
||||||
listSessions(cursor?: string, limit?: number): Promise<AgentSessionList>;
|
|
||||||
getSession(sessionId: string): Promise<AgentSession | null>;
|
|
||||||
getMessages(sessionId: string, limit?: number, before?: string): Promise<AgentMessage[]>;
|
|
||||||
|
|
||||||
// Control operations
|
|
||||||
injectMessage(sessionId: string, content: string): Promise<InjectResult>;
|
|
||||||
pauseSession(sessionId: string): Promise<void>;
|
|
||||||
resumeSession(sessionId: string): Promise<void>;
|
|
||||||
killSession(sessionId: string, force?: boolean): Promise<void>;
|
|
||||||
|
|
||||||
// SSE streaming — returns an AsyncIterable of AgentMessage events
|
|
||||||
streamMessages(sessionId: string): AsyncIterable<AgentMessage>;
|
|
||||||
|
|
||||||
// Health
|
|
||||||
isAvailable(): Promise<boolean>;
|
|
||||||
}
|
|
||||||
|
|
||||||
// Provider configuration stored in DB (AgentProviderConfig model from P0 schema)
|
|
||||||
export interface AgentProviderConfig {
|
|
||||||
id: string;
|
|
||||||
providerId: string;
|
|
||||||
providerType: string;
|
|
||||||
displayName: string;
|
|
||||||
baseUrl?: string;
|
|
||||||
apiToken?: string;
|
|
||||||
enabled: boolean;
|
|
||||||
createdAt: Date;
|
|
||||||
updatedAt: Date;
|
|
||||||
}
|
|
||||||
@@ -134,6 +134,3 @@ export * from "./widget.types";
|
|||||||
|
|
||||||
// Export WebSocket types
|
// Export WebSocket types
|
||||||
export * from "./websocket.types";
|
export * from "./websocket.types";
|
||||||
|
|
||||||
// Export agent provider types
|
|
||||||
export * from "./agent-provider.types";
|
|
||||||
|
|||||||
Reference in New Issue
Block a user