Compare commits
13 Commits
feat/ms23-
...
chore/ms23
| Author | SHA1 | Date | |
|---|---|---|---|
| e57b84af1a | |||
| d0c6622de5 | |||
| 4749f52668 | |||
| e0b28c91c3 | |||
| fa7837af3e | |||
| 123cbce5cd | |||
| 7d47e5ff99 | |||
| ef674206e7 | |||
| 977747599f | |||
| fc4699ca51 | |||
| b61554800b | |||
| 98e892f23c | |||
| de6faf659e |
@@ -21,6 +21,7 @@ FROM base AS deps
|
||||
COPY packages/shared/package.json ./packages/shared/
|
||||
COPY packages/config/package.json ./packages/config/
|
||||
COPY apps/orchestrator/package.json ./apps/orchestrator/
|
||||
# API schema is available via apps/orchestrator/prisma/schema.prisma symlink
|
||||
|
||||
# Copy npm configuration for native binary architecture hints
|
||||
COPY .npmrc ./
|
||||
@@ -46,6 +47,15 @@ COPY --from=deps /app/packages/shared/node_modules ./packages/shared/node_module
|
||||
COPY --from=deps /app/packages/config/node_modules ./packages/config/node_modules
|
||||
COPY --from=deps /app/apps/orchestrator/node_modules ./apps/orchestrator/node_modules
|
||||
|
||||
# The repo has apps/orchestrator/prisma/schema.prisma as a symlink for CI use.
|
||||
# Kaniko resolves destination symlinks on COPY, which fails because the symlink
|
||||
# target (../../api/prisma/schema.prisma) doesn't exist in the container.
|
||||
# Fix: remove the dangling symlink first, then copy the real schema file there.
|
||||
RUN rm -f apps/orchestrator/prisma/schema.prisma
|
||||
COPY apps/api/prisma/schema.prisma ./apps/orchestrator/prisma/schema.prisma
|
||||
# pnpm turbo build runs prisma:generate (--schema=./prisma/schema.prisma) from the
|
||||
# orchestrator package context — no cross-package project-root issues.
|
||||
|
||||
# Build the orchestrator app using TurboRepo
|
||||
RUN pnpm turbo build --filter=@mosaic/orchestrator
|
||||
|
||||
|
||||
@@ -3,19 +3,20 @@
|
||||
"version": "0.0.20",
|
||||
"private": true,
|
||||
"scripts": {
|
||||
"dev": "nest start --watch",
|
||||
"build": "nest build",
|
||||
"dev": "nest start --watch",
|
||||
"lint": "eslint src/",
|
||||
"lint:fix": "eslint src/ --fix",
|
||||
"prisma:generate": "prisma generate --schema=./prisma/schema.prisma",
|
||||
"start": "node dist/main.js",
|
||||
"start:dev": "nest start --watch",
|
||||
"start:debug": "nest start --debug --watch",
|
||||
"start:dev": "nest start --watch",
|
||||
"start:prod": "node dist/main.js",
|
||||
"test": "vitest",
|
||||
"test:watch": "vitest watch",
|
||||
"test:e2e": "vitest run --config tests/integration/vitest.config.ts",
|
||||
"test:perf": "vitest run --config tests/performance/vitest.config.ts",
|
||||
"typecheck": "tsc --noEmit",
|
||||
"lint": "eslint src/",
|
||||
"lint:fix": "eslint src/ --fix"
|
||||
"test:watch": "vitest watch",
|
||||
"typecheck": "tsc --noEmit"
|
||||
},
|
||||
"dependencies": {
|
||||
"@anthropic-ai/sdk": "^0.72.1",
|
||||
@@ -46,6 +47,7 @@
|
||||
"@types/express": "^5.0.1",
|
||||
"@types/node": "^22.13.4",
|
||||
"@vitest/coverage-v8": "^4.0.18",
|
||||
"prisma": "^6.19.2",
|
||||
"ts-node": "^10.9.2",
|
||||
"tsconfig-paths": "^4.2.0",
|
||||
"typescript": "^5.8.2",
|
||||
|
||||
1
apps/orchestrator/prisma/schema.prisma
Symbolic link
1
apps/orchestrator/prisma/schema.prisma
Symbolic link
@@ -0,0 +1 @@
|
||||
../../api/prisma/schema.prisma
|
||||
68
apps/orchestrator/src/api/agents/agent-control.service.ts
Normal file
68
apps/orchestrator/src/api/agents/agent-control.service.ts
Normal file
@@ -0,0 +1,68 @@
|
||||
import { Injectable } from "@nestjs/common";
|
||||
import type { Prisma } from "@prisma/client";
|
||||
import { PrismaService } from "../../prisma/prisma.service";
|
||||
|
||||
@Injectable()
|
||||
export class AgentControlService {
|
||||
constructor(private readonly prisma: PrismaService) {}
|
||||
|
||||
private toJsonValue(value: Record<string, unknown>): Prisma.InputJsonValue {
|
||||
return value as Prisma.InputJsonValue;
|
||||
}
|
||||
|
||||
private async createOperatorAuditLog(
|
||||
agentId: string,
|
||||
operatorId: string,
|
||||
action: "inject" | "pause" | "resume",
|
||||
payload: Record<string, unknown>
|
||||
): Promise<void> {
|
||||
await this.prisma.operatorAuditLog.create({
|
||||
data: {
|
||||
sessionId: agentId,
|
||||
userId: operatorId,
|
||||
provider: "internal",
|
||||
action,
|
||||
metadata: this.toJsonValue({ payload }),
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
async injectMessage(agentId: string, operatorId: string, message: string): Promise<void> {
|
||||
const treeEntry = await this.prisma.agentSessionTree.findUnique({
|
||||
where: { sessionId: agentId },
|
||||
select: { id: true },
|
||||
});
|
||||
|
||||
if (treeEntry) {
|
||||
await this.prisma.agentConversationMessage.create({
|
||||
data: {
|
||||
sessionId: agentId,
|
||||
role: "operator",
|
||||
content: message,
|
||||
provider: "internal",
|
||||
metadata: this.toJsonValue({}),
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
await this.createOperatorAuditLog(agentId, operatorId, "inject", { message });
|
||||
}
|
||||
|
||||
async pauseAgent(agentId: string, operatorId: string): Promise<void> {
|
||||
await this.prisma.agentSessionTree.updateMany({
|
||||
where: { sessionId: agentId },
|
||||
data: { status: "paused" },
|
||||
});
|
||||
|
||||
await this.createOperatorAuditLog(agentId, operatorId, "pause", {});
|
||||
}
|
||||
|
||||
async resumeAgent(agentId: string, operatorId: string): Promise<void> {
|
||||
await this.prisma.agentSessionTree.updateMany({
|
||||
where: { sessionId: agentId },
|
||||
data: { status: "running" },
|
||||
});
|
||||
|
||||
await this.createOperatorAuditLog(agentId, operatorId, "resume", {});
|
||||
}
|
||||
}
|
||||
84
apps/orchestrator/src/api/agents/agent-messages.service.ts
Normal file
84
apps/orchestrator/src/api/agents/agent-messages.service.ts
Normal file
@@ -0,0 +1,84 @@
|
||||
import { Injectable } from "@nestjs/common";
|
||||
import { type AgentConversationMessage, type Prisma } from "@prisma/client";
|
||||
import { PrismaService } from "../../prisma/prisma.service";
|
||||
|
||||
@Injectable()
|
||||
export class AgentMessagesService {
|
||||
constructor(private readonly prisma: PrismaService) {}
|
||||
|
||||
async getMessages(
|
||||
sessionId: string,
|
||||
limit: number,
|
||||
skip: number
|
||||
): Promise<{
|
||||
messages: AgentConversationMessage[];
|
||||
total: number;
|
||||
}> {
|
||||
const where = { sessionId };
|
||||
|
||||
const [messages, total] = await Promise.all([
|
||||
this.prisma.agentConversationMessage.findMany({
|
||||
where,
|
||||
orderBy: {
|
||||
timestamp: "desc",
|
||||
},
|
||||
take: limit,
|
||||
skip,
|
||||
}),
|
||||
this.prisma.agentConversationMessage.count({ where }),
|
||||
]);
|
||||
|
||||
return {
|
||||
messages,
|
||||
total,
|
||||
};
|
||||
}
|
||||
|
||||
async getReplayMessages(sessionId: string, limit = 50): Promise<AgentConversationMessage[]> {
|
||||
const messages = await this.prisma.agentConversationMessage.findMany({
|
||||
where: { sessionId },
|
||||
orderBy: {
|
||||
timestamp: "desc",
|
||||
},
|
||||
take: limit,
|
||||
});
|
||||
|
||||
return messages.reverse();
|
||||
}
|
||||
|
||||
async getMessagesAfter(
|
||||
sessionId: string,
|
||||
lastSeenTimestamp: Date,
|
||||
lastSeenMessageId: string | null
|
||||
): Promise<AgentConversationMessage[]> {
|
||||
const where: Prisma.AgentConversationMessageWhereInput = {
|
||||
sessionId,
|
||||
...(lastSeenMessageId
|
||||
? {
|
||||
OR: [
|
||||
{
|
||||
timestamp: {
|
||||
gt: lastSeenTimestamp,
|
||||
},
|
||||
},
|
||||
{
|
||||
timestamp: lastSeenTimestamp,
|
||||
id: {
|
||||
gt: lastSeenMessageId,
|
||||
},
|
||||
},
|
||||
],
|
||||
}
|
||||
: {
|
||||
timestamp: {
|
||||
gt: lastSeenTimestamp,
|
||||
},
|
||||
}),
|
||||
};
|
||||
|
||||
return this.prisma.agentConversationMessage.findMany({
|
||||
where,
|
||||
orderBy: [{ timestamp: "asc" }, { id: "asc" }],
|
||||
});
|
||||
}
|
||||
}
|
||||
@@ -5,6 +5,8 @@ import { AgentSpawnerService } from "../../spawner/agent-spawner.service";
|
||||
import { AgentLifecycleService } from "../../spawner/agent-lifecycle.service";
|
||||
import { KillswitchService } from "../../killswitch/killswitch.service";
|
||||
import { AgentEventsService } from "./agent-events.service";
|
||||
import { AgentMessagesService } from "./agent-messages.service";
|
||||
import { AgentControlService } from "./agent-control.service";
|
||||
import type { KillAllResult } from "../../killswitch/killswitch.service";
|
||||
|
||||
describe("AgentsController - Killswitch Endpoints", () => {
|
||||
@@ -27,6 +29,17 @@ describe("AgentsController - Killswitch Endpoints", () => {
|
||||
subscribe: ReturnType<typeof vi.fn>;
|
||||
getInitialSnapshot: ReturnType<typeof vi.fn>;
|
||||
createHeartbeat: ReturnType<typeof vi.fn>;
|
||||
getRecentEvents: ReturnType<typeof vi.fn>;
|
||||
};
|
||||
let mockMessagesService: {
|
||||
getMessages: ReturnType<typeof vi.fn>;
|
||||
getReplayMessages: ReturnType<typeof vi.fn>;
|
||||
getMessagesAfter: ReturnType<typeof vi.fn>;
|
||||
};
|
||||
let mockControlService: {
|
||||
injectMessage: ReturnType<typeof vi.fn>;
|
||||
pauseAgent: ReturnType<typeof vi.fn>;
|
||||
resumeAgent: ReturnType<typeof vi.fn>;
|
||||
};
|
||||
|
||||
beforeEach(() => {
|
||||
@@ -61,6 +74,19 @@ describe("AgentsController - Killswitch Endpoints", () => {
|
||||
timestamp: new Date().toISOString(),
|
||||
data: { heartbeat: true },
|
||||
}),
|
||||
getRecentEvents: vi.fn().mockReturnValue([]),
|
||||
};
|
||||
|
||||
mockMessagesService = {
|
||||
getMessages: vi.fn(),
|
||||
getReplayMessages: vi.fn().mockResolvedValue([]),
|
||||
getMessagesAfter: vi.fn().mockResolvedValue([]),
|
||||
};
|
||||
|
||||
mockControlService = {
|
||||
injectMessage: vi.fn().mockResolvedValue(undefined),
|
||||
pauseAgent: vi.fn().mockResolvedValue(undefined),
|
||||
resumeAgent: vi.fn().mockResolvedValue(undefined),
|
||||
};
|
||||
|
||||
controller = new AgentsController(
|
||||
@@ -68,7 +94,9 @@ describe("AgentsController - Killswitch Endpoints", () => {
|
||||
mockSpawnerService as unknown as AgentSpawnerService,
|
||||
mockLifecycleService as unknown as AgentLifecycleService,
|
||||
mockKillswitchService as unknown as KillswitchService,
|
||||
mockEventsService as unknown as AgentEventsService
|
||||
mockEventsService as unknown as AgentEventsService,
|
||||
mockMessagesService as unknown as AgentMessagesService,
|
||||
mockControlService as unknown as AgentControlService
|
||||
);
|
||||
});
|
||||
|
||||
|
||||
@@ -4,6 +4,8 @@ import { AgentSpawnerService } from "../../spawner/agent-spawner.service";
|
||||
import { AgentLifecycleService } from "../../spawner/agent-lifecycle.service";
|
||||
import { KillswitchService } from "../../killswitch/killswitch.service";
|
||||
import { AgentEventsService } from "./agent-events.service";
|
||||
import { AgentMessagesService } from "./agent-messages.service";
|
||||
import { AgentControlService } from "./agent-control.service";
|
||||
import { describe, it, expect, beforeEach, afterEach, vi } from "vitest";
|
||||
|
||||
describe("AgentsController", () => {
|
||||
@@ -30,6 +32,16 @@ describe("AgentsController", () => {
|
||||
createHeartbeat: ReturnType<typeof vi.fn>;
|
||||
getRecentEvents: ReturnType<typeof vi.fn>;
|
||||
};
|
||||
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>;
|
||||
};
|
||||
|
||||
beforeEach(() => {
|
||||
// Create mock services
|
||||
@@ -69,13 +81,27 @@ describe("AgentsController", () => {
|
||||
getRecentEvents: vi.fn().mockReturnValue([]),
|
||||
};
|
||||
|
||||
messagesService = {
|
||||
getMessages: vi.fn(),
|
||||
getReplayMessages: vi.fn().mockResolvedValue([]),
|
||||
getMessagesAfter: vi.fn().mockResolvedValue([]),
|
||||
};
|
||||
|
||||
controlService = {
|
||||
injectMessage: vi.fn().mockResolvedValue(undefined),
|
||||
pauseAgent: vi.fn().mockResolvedValue(undefined),
|
||||
resumeAgent: vi.fn().mockResolvedValue(undefined),
|
||||
};
|
||||
|
||||
// Create controller with mocked services
|
||||
controller = new AgentsController(
|
||||
queueService as unknown as QueueService,
|
||||
spawnerService as unknown as AgentSpawnerService,
|
||||
lifecycleService as unknown as AgentLifecycleService,
|
||||
killswitchService as unknown as KillswitchService,
|
||||
eventsService as unknown as AgentEventsService
|
||||
eventsService as unknown as AgentEventsService,
|
||||
messagesService as unknown as AgentMessagesService,
|
||||
controlService as unknown as AgentControlService
|
||||
);
|
||||
});
|
||||
|
||||
@@ -365,6 +391,93 @@ describe("AgentsController", () => {
|
||||
});
|
||||
});
|
||||
|
||||
describe("agent control endpoints", () => {
|
||||
const agentId = "0b64079f-4487-42b9-92eb-cf8ea0042a64";
|
||||
|
||||
it("should inject an operator message", async () => {
|
||||
const req = { apiKey: "control-key" };
|
||||
|
||||
const result = await controller.injectAgentMessage(
|
||||
agentId,
|
||||
{ message: "pause and summarize" },
|
||||
req
|
||||
);
|
||||
|
||||
expect(controlService.injectMessage).toHaveBeenCalledWith(
|
||||
agentId,
|
||||
"control-key",
|
||||
"pause and summarize"
|
||||
);
|
||||
expect(result).toEqual({ message: `Message injected into agent ${agentId}` });
|
||||
});
|
||||
|
||||
it("should default operator id when request api key is missing", async () => {
|
||||
await controller.injectAgentMessage(agentId, { message: "continue" }, {});
|
||||
|
||||
expect(controlService.injectMessage).toHaveBeenCalledWith(agentId, "operator", "continue");
|
||||
});
|
||||
|
||||
it("should pause an agent", async () => {
|
||||
const result = await controller.pauseAgent(agentId, {}, { apiKey: "ops-user" });
|
||||
|
||||
expect(controlService.pauseAgent).toHaveBeenCalledWith(agentId, "ops-user");
|
||||
expect(result).toEqual({ message: `Agent ${agentId} paused` });
|
||||
});
|
||||
|
||||
it("should resume an agent", async () => {
|
||||
const result = await controller.resumeAgent(agentId, {}, { apiKey: "ops-user" });
|
||||
|
||||
expect(controlService.resumeAgent).toHaveBeenCalledWith(agentId, "ops-user");
|
||||
expect(result).toEqual({ message: `Agent ${agentId} resumed` });
|
||||
});
|
||||
});
|
||||
|
||||
describe("getAgentMessages", () => {
|
||||
it("should return paginated message history", async () => {
|
||||
const agentId = "0b64079f-4487-42b9-92eb-cf8ea0042a64";
|
||||
const query = {
|
||||
limit: 25,
|
||||
skip: 10,
|
||||
};
|
||||
|
||||
const response = {
|
||||
messages: [
|
||||
{
|
||||
id: "msg-1",
|
||||
sessionId: agentId,
|
||||
role: "agent",
|
||||
content: "hello",
|
||||
provider: "internal",
|
||||
timestamp: new Date("2026-03-07T03:00:00.000Z"),
|
||||
metadata: {},
|
||||
},
|
||||
],
|
||||
total: 101,
|
||||
};
|
||||
|
||||
messagesService.getMessages.mockResolvedValue(response);
|
||||
|
||||
const result = await controller.getAgentMessages(agentId, query);
|
||||
|
||||
expect(messagesService.getMessages).toHaveBeenCalledWith(agentId, 25, 10);
|
||||
expect(result).toEqual(response);
|
||||
});
|
||||
|
||||
it("should use default pagination values", async () => {
|
||||
const agentId = "0b64079f-4487-42b9-92eb-cf8ea0042a64";
|
||||
const query = {
|
||||
limit: 50,
|
||||
skip: 0,
|
||||
};
|
||||
|
||||
messagesService.getMessages.mockResolvedValue({ messages: [], total: 0 });
|
||||
|
||||
await controller.getAgentMessages(agentId, query);
|
||||
|
||||
expect(messagesService.getMessages).toHaveBeenCalledWith(agentId, 50, 0);
|
||||
});
|
||||
});
|
||||
|
||||
describe("getRecentEvents", () => {
|
||||
it("should return recent events with default limit", () => {
|
||||
eventsService.getRecentEvents.mockReturnValue([
|
||||
|
||||
@@ -14,7 +14,9 @@ import {
|
||||
Sse,
|
||||
MessageEvent,
|
||||
Query,
|
||||
Request,
|
||||
} from "@nestjs/common";
|
||||
import type { AgentConversationMessage } from "@prisma/client";
|
||||
import { Throttle } from "@nestjs/throttler";
|
||||
import { Observable } from "rxjs";
|
||||
import { QueueService } from "../../queue/queue.service";
|
||||
@@ -25,6 +27,11 @@ import { SpawnAgentDto, SpawnAgentResponseDto } from "./dto/spawn-agent.dto";
|
||||
import { OrchestratorApiKeyGuard } from "../../common/guards/api-key.guard";
|
||||
import { OrchestratorThrottlerGuard } from "../../common/guards/throttler.guard";
|
||||
import { AgentEventsService } from "./agent-events.service";
|
||||
import { GetMessagesQueryDto } from "./dto/get-messages-query.dto";
|
||||
import { AgentMessagesService } from "./agent-messages.service";
|
||||
import { AgentControlService } from "./agent-control.service";
|
||||
import { InjectAgentDto } from "./dto/inject-agent.dto";
|
||||
import { PauseAgentDto, ResumeAgentDto } from "./dto/control-agent.dto";
|
||||
|
||||
/**
|
||||
* Controller for agent management endpoints
|
||||
@@ -47,7 +54,9 @@ export class AgentsController {
|
||||
private readonly spawnerService: AgentSpawnerService,
|
||||
private readonly lifecycleService: AgentLifecycleService,
|
||||
private readonly killswitchService: KillswitchService,
|
||||
private readonly eventsService: AgentEventsService
|
||||
private readonly eventsService: AgentEventsService,
|
||||
private readonly messagesService: AgentMessagesService,
|
||||
private readonly agentControlService: AgentControlService
|
||||
) {}
|
||||
|
||||
/**
|
||||
@@ -185,6 +194,107 @@ export class AgentsController {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Get paginated message history for an agent.
|
||||
*/
|
||||
@Get(":agentId/messages")
|
||||
@Throttle({ status: { limit: 200, ttl: 60000 } })
|
||||
@UsePipes(new ValidationPipe({ transform: true, whitelist: true }))
|
||||
async getAgentMessages(
|
||||
@Param("agentId", ParseUUIDPipe) agentId: string,
|
||||
@Query() query: GetMessagesQueryDto
|
||||
): Promise<{
|
||||
messages: AgentConversationMessage[];
|
||||
total: number;
|
||||
}> {
|
||||
return this.messagesService.getMessages(agentId, query.limit, query.skip);
|
||||
}
|
||||
|
||||
/**
|
||||
* Stream per-agent conversation messages as server-sent events (SSE).
|
||||
*/
|
||||
@Sse(":agentId/messages/stream")
|
||||
@Throttle({ status: { limit: 200, ttl: 60000 } })
|
||||
streamAgentMessages(@Param("agentId", ParseUUIDPipe) agentId: string): Observable<MessageEvent> {
|
||||
return new Observable<MessageEvent>((subscriber) => {
|
||||
let isClosed = false;
|
||||
let lastSeenTimestamp = new Date();
|
||||
let lastSeenMessageId: string | null = null;
|
||||
|
||||
const emitMessage = (message: AgentConversationMessage): void => {
|
||||
if (isClosed) {
|
||||
return;
|
||||
}
|
||||
|
||||
subscriber.next({
|
||||
data: this.toMessageStreamPayload(message),
|
||||
});
|
||||
|
||||
lastSeenTimestamp = message.timestamp;
|
||||
lastSeenMessageId = message.id;
|
||||
};
|
||||
|
||||
void this.messagesService
|
||||
.getReplayMessages(agentId, 50)
|
||||
.then((messages) => {
|
||||
if (isClosed) {
|
||||
return;
|
||||
}
|
||||
|
||||
messages.forEach((message) => {
|
||||
emitMessage(message);
|
||||
});
|
||||
|
||||
if (messages.length === 0) {
|
||||
lastSeenTimestamp = new Date();
|
||||
lastSeenMessageId = null;
|
||||
}
|
||||
})
|
||||
.catch((error: unknown) => {
|
||||
this.logger.error(
|
||||
`Failed to load replay messages for ${agentId}: ${error instanceof Error ? error.message : String(error)}`
|
||||
);
|
||||
lastSeenTimestamp = new Date();
|
||||
lastSeenMessageId = null;
|
||||
});
|
||||
|
||||
const pollInterval = setInterval(() => {
|
||||
if (isClosed) {
|
||||
return;
|
||||
}
|
||||
|
||||
void this.messagesService
|
||||
.getMessagesAfter(agentId, lastSeenTimestamp, lastSeenMessageId)
|
||||
.then((messages) => {
|
||||
if (isClosed || messages.length === 0) {
|
||||
return;
|
||||
}
|
||||
|
||||
messages.forEach((message) => {
|
||||
emitMessage(message);
|
||||
});
|
||||
})
|
||||
.catch((error: unknown) => {
|
||||
this.logger.error(
|
||||
`Failed to poll messages for ${agentId}: ${error instanceof Error ? error.message : String(error)}`
|
||||
);
|
||||
});
|
||||
}, 1000);
|
||||
|
||||
const heartbeat = setInterval(() => {
|
||||
if (!isClosed) {
|
||||
subscriber.next({ data: { type: "heartbeat" } });
|
||||
}
|
||||
}, 15000);
|
||||
|
||||
return () => {
|
||||
isClosed = true;
|
||||
clearInterval(pollInterval);
|
||||
clearInterval(heartbeat);
|
||||
};
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* Get agent status
|
||||
* @param agentId Agent ID to query
|
||||
@@ -269,6 +379,57 @@ export class AgentsController {
|
||||
}
|
||||
}
|
||||
|
||||
@Post(":agentId/inject")
|
||||
@Throttle({ default: { limit: 10, ttl: 60000 } })
|
||||
@HttpCode(200)
|
||||
@UsePipes(new ValidationPipe({ transform: true, whitelist: true }))
|
||||
async injectAgentMessage(
|
||||
@Param("agentId", ParseUUIDPipe) agentId: string,
|
||||
@Body() dto: InjectAgentDto,
|
||||
@Request() req: { apiKey?: string }
|
||||
): Promise<{ message: string }> {
|
||||
const operatorId = req.apiKey ?? "operator";
|
||||
await this.agentControlService.injectMessage(agentId, operatorId, dto.message);
|
||||
|
||||
return {
|
||||
message: `Message injected into agent ${agentId}`,
|
||||
};
|
||||
}
|
||||
|
||||
@Post(":agentId/pause")
|
||||
@Throttle({ default: { limit: 10, ttl: 60000 } })
|
||||
@HttpCode(200)
|
||||
@UsePipes(new ValidationPipe({ transform: true, whitelist: true }))
|
||||
async pauseAgent(
|
||||
@Param("agentId", ParseUUIDPipe) agentId: string,
|
||||
@Body() _dto: PauseAgentDto,
|
||||
@Request() req: { apiKey?: string }
|
||||
): Promise<{ message: string }> {
|
||||
const operatorId = req.apiKey ?? "operator";
|
||||
await this.agentControlService.pauseAgent(agentId, operatorId);
|
||||
|
||||
return {
|
||||
message: `Agent ${agentId} paused`,
|
||||
};
|
||||
}
|
||||
|
||||
@Post(":agentId/resume")
|
||||
@Throttle({ default: { limit: 10, ttl: 60000 } })
|
||||
@HttpCode(200)
|
||||
@UsePipes(new ValidationPipe({ transform: true, whitelist: true }))
|
||||
async resumeAgent(
|
||||
@Param("agentId", ParseUUIDPipe) agentId: string,
|
||||
@Body() _dto: ResumeAgentDto,
|
||||
@Request() req: { apiKey?: string }
|
||||
): Promise<{ message: string }> {
|
||||
const operatorId = req.apiKey ?? "operator";
|
||||
await this.agentControlService.resumeAgent(agentId, operatorId);
|
||||
|
||||
return {
|
||||
message: `Agent ${agentId} resumed`,
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* Kill all active agents
|
||||
* @returns Summary of kill operation
|
||||
@@ -301,4 +462,24 @@ export class AgentsController {
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
private toMessageStreamPayload(message: AgentConversationMessage): {
|
||||
messageId: string;
|
||||
sessionId: string;
|
||||
role: string;
|
||||
content: string;
|
||||
provider: string;
|
||||
timestamp: string;
|
||||
metadata: unknown;
|
||||
} {
|
||||
return {
|
||||
messageId: message.id,
|
||||
sessionId: message.sessionId,
|
||||
role: message.role,
|
||||
content: message.content,
|
||||
provider: message.provider,
|
||||
timestamp: message.timestamp.toISOString(),
|
||||
metadata: message.metadata,
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
@@ -6,10 +6,18 @@ import { KillswitchModule } from "../../killswitch/killswitch.module";
|
||||
import { ValkeyModule } from "../../valkey/valkey.module";
|
||||
import { OrchestratorApiKeyGuard } from "../../common/guards/api-key.guard";
|
||||
import { AgentEventsService } from "./agent-events.service";
|
||||
import { PrismaModule } from "../../prisma/prisma.module";
|
||||
import { AgentMessagesService } from "./agent-messages.service";
|
||||
import { AgentControlService } from "./agent-control.service";
|
||||
|
||||
@Module({
|
||||
imports: [QueueModule, SpawnerModule, KillswitchModule, ValkeyModule],
|
||||
imports: [QueueModule, SpawnerModule, KillswitchModule, ValkeyModule, PrismaModule],
|
||||
controllers: [AgentsController],
|
||||
providers: [OrchestratorApiKeyGuard, AgentEventsService],
|
||||
providers: [
|
||||
OrchestratorApiKeyGuard,
|
||||
AgentEventsService,
|
||||
AgentMessagesService,
|
||||
AgentControlService,
|
||||
],
|
||||
})
|
||||
export class AgentsModule {}
|
||||
|
||||
@@ -0,0 +1,3 @@
|
||||
export class PauseAgentDto {}
|
||||
|
||||
export class ResumeAgentDto {}
|
||||
@@ -0,0 +1,37 @@
|
||||
import { plainToInstance } from "class-transformer";
|
||||
import { validate } from "class-validator";
|
||||
import { describe, expect, it } from "vitest";
|
||||
import { GetMessagesQueryDto } from "./get-messages-query.dto";
|
||||
|
||||
describe("GetMessagesQueryDto", () => {
|
||||
it("should use defaults when empty", async () => {
|
||||
const dto = plainToInstance(GetMessagesQueryDto, {});
|
||||
const errors = await validate(dto);
|
||||
|
||||
expect(errors).toHaveLength(0);
|
||||
expect(dto.limit).toBe(50);
|
||||
expect(dto.skip).toBe(0);
|
||||
});
|
||||
|
||||
it("should reject limit greater than 200", async () => {
|
||||
const dto = plainToInstance(GetMessagesQueryDto, {
|
||||
limit: 201,
|
||||
skip: 0,
|
||||
});
|
||||
const errors = await validate(dto);
|
||||
|
||||
expect(errors.length).toBeGreaterThan(0);
|
||||
expect(errors.some((error) => error.property === "limit")).toBe(true);
|
||||
});
|
||||
|
||||
it("should reject negative skip", async () => {
|
||||
const dto = plainToInstance(GetMessagesQueryDto, {
|
||||
limit: 50,
|
||||
skip: -1,
|
||||
});
|
||||
const errors = await validate(dto);
|
||||
|
||||
expect(errors.length).toBeGreaterThan(0);
|
||||
expect(errors.some((error) => error.property === "skip")).toBe(true);
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,17 @@
|
||||
import { Type } from "class-transformer";
|
||||
import { IsInt, IsOptional, Max, Min } from "class-validator";
|
||||
|
||||
export class GetMessagesQueryDto {
|
||||
@IsOptional()
|
||||
@Type(() => Number)
|
||||
@IsInt()
|
||||
@Min(1)
|
||||
@Max(200)
|
||||
limit = 50;
|
||||
|
||||
@IsOptional()
|
||||
@Type(() => Number)
|
||||
@IsInt()
|
||||
@Min(0)
|
||||
skip = 0;
|
||||
}
|
||||
7
apps/orchestrator/src/api/agents/dto/inject-agent.dto.ts
Normal file
7
apps/orchestrator/src/api/agents/dto/inject-agent.dto.ts
Normal file
@@ -0,0 +1,7 @@
|
||||
import { IsNotEmpty, IsString } from "class-validator";
|
||||
|
||||
export class InjectAgentDto {
|
||||
@IsString()
|
||||
@IsNotEmpty()
|
||||
message!: string;
|
||||
}
|
||||
@@ -4,6 +4,7 @@ export default defineConfig({
|
||||
test: {
|
||||
globals: true,
|
||||
environment: "node",
|
||||
setupFiles: ["reflect-metadata"],
|
||||
exclude: ["**/node_modules/**", "**/dist/**", "**/tests/integration/**"],
|
||||
include: ["src/**/*.spec.ts", "src/**/*.test.ts"],
|
||||
coverage: {
|
||||
|
||||
@@ -122,10 +122,10 @@ Target version: `v0.0.23`
|
||||
| id | status | milestone | description | issue | repo | branch | depends_on | blocks | agent | started_at | completed_at | estimate | used | notes |
|
||||
| ----------- | ----------- | ------------- | ------------------------------------------------------------------------------------------------ | ----- | ------------ | ---------------------- | ----------------------------------------------- | ----------------------------------------------------------- | ----- | ---------- | ------------ | -------- | ---- | --------------------------------------------- |
|
||||
| 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 | in-progress | 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 | — | 20K | — | |
|
||||
| MS23-P0-003 | not-started | p0-foundation | Orchestrator API: GET /agents/:id/messages + SSE stream endpoint | #693 | orchestrator | feat/ms23-p0-stream | MS23-P0-001 | MS23-P0-006 | — | — | — | 20K | — | |
|
||||
| MS23-P0-004 | not-started | p0-foundation | Orchestrator API: POST /agents/:id/inject + pause/resume endpoints | #693 | orchestrator | feat/ms23-p0-controls | 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-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-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-005 | in-progress | 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 | 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)
|
||||
|
||||
3
pnpm-lock.yaml
generated
3
pnpm-lock.yaml
generated
@@ -389,6 +389,9 @@ importers:
|
||||
'@vitest/coverage-v8':
|
||||
specifier: ^4.0.18
|
||||
version: 4.0.18(vitest@4.0.18(@opentelemetry/api@1.9.0)(@types/node@22.19.7)(jiti@2.6.1)(jsdom@26.1.0)(terser@5.46.0)(tsx@4.21.0)(yaml@2.8.2))
|
||||
prisma:
|
||||
specifier: ^6.19.2
|
||||
version: 6.19.2(magicast@0.3.5)(typescript@5.9.3)
|
||||
ts-node:
|
||||
specifier: ^10.9.2
|
||||
version: 10.9.2(@swc/core@1.15.11)(@types/node@22.19.7)(typescript@5.9.3)
|
||||
|
||||
Reference in New Issue
Block a user