ci/woodpecker/push/publish Pipeline failed
Co-authored-by: shaggy <[email protected]>
959 lines
41 KiB
TypeScript
959 lines
41 KiB
TypeScript
import 'reflect-metadata';
|
|
import { readFileSync } from 'node:fs';
|
|
import { resolve } from 'node:path';
|
|
import { ForbiddenException, NotFoundException } from '@nestjs/common';
|
|
import { Test, type TestingModule } from '@nestjs/testing';
|
|
import { describe, expect, it, vi } from 'vitest';
|
|
|
|
vi.mock('../agent.service.js', () => ({ AgentService: class AgentService {} }));
|
|
vi.mock('../../commands/command-executor.service.js', () => ({
|
|
CommandExecutorService: class CommandExecutorService {},
|
|
}));
|
|
vi.mock('../routing/routing-engine.service.js', () => ({
|
|
RoutingEngineService: class RoutingEngineService {},
|
|
}));
|
|
|
|
import { SessionsController } from '../sessions.controller.js';
|
|
import { AgentService } from '../agent.service.js';
|
|
import { ChatController } from '../../chat/chat.controller.js';
|
|
import { ChatGateway } from '../../chat/chat.gateway.js';
|
|
import type { AgentSession } from '../agent.service.js';
|
|
import type { SessionInfoDto } from '../session.dto.js';
|
|
import type { HarnessAdapter, HarnessConversationService } from '@mosaicstack/types';
|
|
import { AuthGuard } from '../../auth/auth.guard.js';
|
|
import { AUTH } from '../../auth/auth.tokens.js';
|
|
import { BRAIN } from '../../brain/brain.tokens.js';
|
|
import { CommandRegistryService } from '../../commands/command-registry.service.js';
|
|
import { CommandExecutorService } from '../../commands/command-executor.service.js';
|
|
import { RoutingEngineService } from '../routing/routing-engine.service.js';
|
|
import { ChatRuntimeRouter } from '../../chat/chat-runtime-router.js';
|
|
import { EmbeddedChatRuntime } from '../../chat/embedded-chat.runtime.js';
|
|
import { ownConversation } from '../../chat/chat-runtime.js';
|
|
import type { LegacyRuntimeStream } from '../../chat/chat-runtime.js';
|
|
import { HarnessChatRuntime } from '../../chat/harness-chat.runtime.js';
|
|
import { HarnessRegistry } from '../../harness/harness.registry.js';
|
|
import { HARNESS_CONVERSATION_SERVICE_UNAVAILABLE } from '../../harness/harness.tokens.js';
|
|
|
|
const USER_A = { id: 'user-a', tenantId: 'tenant-a' };
|
|
const USER_B = { id: 'user-b', tenantId: 'tenant-b' };
|
|
const CONVERSATION_ID = '11111111-1111-4111-8111-111111111111';
|
|
|
|
function makeSessionInfo(overrides?: Partial<SessionInfoDto>): SessionInfoDto {
|
|
return {
|
|
id: CONVERSATION_ID,
|
|
provider: 'test-provider',
|
|
modelId: 'test-model',
|
|
createdAt: new Date('2026-07-12T00:00:00Z').toISOString(),
|
|
promptCount: 0,
|
|
channels: [],
|
|
durationMs: 0,
|
|
metrics: {
|
|
tokens: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
|
|
modelSwitches: 0,
|
|
messageCount: 0,
|
|
lastActivityAt: new Date('2026-07-12T00:00:00Z').toISOString(),
|
|
},
|
|
...overrides,
|
|
};
|
|
}
|
|
|
|
function makeAgentSession(owner = USER_A): AgentSession {
|
|
return {
|
|
id: CONVERSATION_ID,
|
|
provider: 'test-provider',
|
|
modelId: 'test-model',
|
|
piSession: {
|
|
thinkingLevel: 'off',
|
|
getAvailableThinkingLevels: vi.fn().mockReturnValue(['off', 'low', 'high']),
|
|
setThinkingLevel: vi.fn(),
|
|
abort: vi.fn().mockResolvedValue(undefined),
|
|
prompt: vi.fn().mockResolvedValue(undefined),
|
|
dispose: vi.fn(),
|
|
getSessionStats: vi.fn(),
|
|
getContextUsage: vi.fn(),
|
|
} as unknown as AgentSession['piSession'],
|
|
listeners: new Set(),
|
|
unsubscribe: vi.fn(),
|
|
createdAt: Date.now(),
|
|
promptCount: 0,
|
|
channels: new Set(),
|
|
skillPromptAdditions: [],
|
|
sandboxDir: '/tmp',
|
|
allowedTools: null,
|
|
userId: owner.id,
|
|
tenantId: owner.tenantId,
|
|
metrics: {
|
|
tokens: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
|
|
modelSwitches: 0,
|
|
messageCount: 0,
|
|
lastActivityAt: new Date('2026-07-12T00:00:00Z').toISOString(),
|
|
},
|
|
};
|
|
}
|
|
|
|
/**
|
|
* A shape-complete, non-throwing AgentService fake scoped so that USER_B (a foreign owner guessing
|
|
* USER_A's conversation id) is never granted the session. Because every method exists and no method
|
|
* throws for a wrong shape, production runs to its real ownership decision — the RED never comes from
|
|
* a `getSession is not a function` TypeError, only from a router-boundary/scope assertion mismatch.
|
|
*/
|
|
function makeScopedAgentService() {
|
|
const foreign = makeAgentSession(USER_A);
|
|
return {
|
|
listSessions: vi.fn((scope?: { userId: string; tenantId?: string }) =>
|
|
scope?.userId === USER_B.id ? [] : [makeSessionInfo({ id: foreign.id })],
|
|
),
|
|
getSessionInfo: vi.fn((_id: string, scope?: { userId: string; tenantId?: string }) =>
|
|
scope?.userId === USER_B.id ? undefined : makeSessionInfo({ id: foreign.id }),
|
|
),
|
|
destroySession: vi.fn(),
|
|
getSession: vi.fn((_id: string, scope?: { userId: string; tenantId?: string }) =>
|
|
scope?.userId === USER_B.id ? undefined : foreign,
|
|
),
|
|
createSession: vi.fn().mockRejectedValue(new NotFoundException('Session scope mismatch')),
|
|
onEvent: vi.fn(() => vi.fn()),
|
|
addChannel: vi.fn(),
|
|
removeChannel: vi.fn(),
|
|
recordMessage: vi.fn(),
|
|
prompt: vi.fn().mockResolvedValue(undefined),
|
|
};
|
|
}
|
|
|
|
type ScopedAgentService = ReturnType<typeof makeScopedAgentService>;
|
|
|
|
/**
|
|
* A structurally-complete harness conversation service that throws if any method is invoked.
|
|
* Fronted behind the legacy runtime's harness slot: the legacy path must never reach it.
|
|
*/
|
|
const failIfUsedConversationService = {
|
|
attach: () => {
|
|
throw new Error('harness conversation service must not be reached on the legacy path');
|
|
},
|
|
detach: () => {
|
|
throw new Error('harness conversation service must not be reached on the legacy path');
|
|
},
|
|
send: () => {
|
|
throw new Error('harness conversation service must not be reached on the legacy path');
|
|
},
|
|
|
|
subscribeFrom: async function* () {
|
|
throw new Error('harness conversation service must not be reached on the legacy path');
|
|
},
|
|
} as unknown as HarnessConversationService;
|
|
|
|
/** A structurally-complete, non-sentinel conversation service used to satisfy the pi-rpc readiness gate. */
|
|
const boundConversationService = {
|
|
attach: () => Promise.reject(new Error('unused')),
|
|
detach: () => Promise.reject(new Error('unused')),
|
|
send: () => Promise.reject(new Error('unused')),
|
|
|
|
subscribeFrom: async function* () {
|
|
throw new Error('unused');
|
|
},
|
|
} as unknown as HarnessConversationService;
|
|
|
|
function registryWith(adapterIds: readonly string[]): HarnessRegistry {
|
|
const registry = new HarnessRegistry();
|
|
for (const id of adapterIds) {
|
|
registry.register({
|
|
id,
|
|
describe: () => Promise.reject(new Error('unused')),
|
|
catalog: () => Promise.reject(new Error('unused')),
|
|
create: () => Promise.reject(new Error('unused')),
|
|
resume: () => Promise.reject(new Error('unused')),
|
|
} as HarnessAdapter);
|
|
}
|
|
return registry;
|
|
}
|
|
|
|
/**
|
|
* Build the real legacy-mode {@link ChatRuntimeRouter} fronting a real {@link EmbeddedChatRuntime}
|
|
* that holds the scoped AgentService fake. This is the ONLY path server-derived scope may travel to
|
|
* reach an AgentService: controller/gateway → ChatRuntimeRouter → EmbeddedChatRuntime → AgentService.
|
|
* The `embeddedAgentService` handed here is a SEPARATE instance from the directly-injected fake, so a
|
|
* call landing on it proves the router-delegation redesign is live rather than the old direct path.
|
|
*/
|
|
function legacyRouterFronting(agentService: unknown): ChatRuntimeRouter {
|
|
const embedded = new EmbeddedChatRuntime(agentService as never);
|
|
const harness = new HarnessChatRuntime(failIfUsedConversationService);
|
|
const router = new ChatRuntimeRouter(
|
|
new HarnessRegistry(),
|
|
HARNESS_CONVERSATION_SERVICE_UNAVAILABLE,
|
|
embedded,
|
|
harness,
|
|
'legacy',
|
|
);
|
|
router.onModuleInit();
|
|
return router;
|
|
}
|
|
|
|
/**
|
|
* The AgentService method names the controller/gateway must NEVER drive on the runtime at the
|
|
* delegation boundary. An AgentService-shaped router shim (a method-for-method mirror) would record
|
|
* one of these instead of the frozen legacy op, so asserting their ABSENCE from the observed runtime
|
|
* call set defeats the shim on INVOCATION evidence — never satisfiable by dead source text.
|
|
*/
|
|
const FORBIDDEN_AGENT_OPS = [
|
|
'getSession',
|
|
'createSession',
|
|
'onEvent',
|
|
'addChannel',
|
|
'prompt',
|
|
'setThinking',
|
|
'abort',
|
|
] as const;
|
|
|
|
/**
|
|
* Wrap a real {@link ChatRuntimeRouter} in a call-recording Proxy. Every property access that yields
|
|
* an OWN/inherited callable is returned as a thin wrapper that appends the method name to `calls` at
|
|
* INVOCATION time and forwards to the real method (bound to the real target, so the router's internal
|
|
* delegation to the embedded runtime runs untouched below this boundary). Non-function and MISSING
|
|
* properties are returned verbatim via Reflect.get — the observer NEVER fabricates a value, returns a
|
|
* canned outcome, or delegates a not-yet-implemented named op, so it cannot itself become a shim.
|
|
*
|
|
* The result is a RUNTIME call set of exactly the methods the controller/gateway invoke ON the router
|
|
* at the delegation seam. Only an actual call can enter it; a dead method, comment, or string in the
|
|
* production source cannot. This replaces the earlier `source.toContain('<frozen op>')` proof — which
|
|
* a dead declaration could satisfy while production still executed a shim — with invocation evidence.
|
|
*/
|
|
function makeRecordingRouter(target: ChatRuntimeRouter, calls: string[]): ChatRuntimeRouter {
|
|
return new Proxy(target, {
|
|
get(t, prop) {
|
|
const value = Reflect.get(t, prop);
|
|
if (typeof value === 'function' && typeof prop === 'string') {
|
|
return (...args: unknown[]) => {
|
|
calls.push(prop);
|
|
return (value as (...a: unknown[]) => unknown).apply(t, args);
|
|
};
|
|
}
|
|
return value;
|
|
},
|
|
}) as ChatRuntimeRouter;
|
|
}
|
|
|
|
/**
|
|
* Real Nest DI dual-provider fixture (mirrors the blessed group-3 pattern in chat-security.test.ts).
|
|
*
|
|
* BOTH an `AgentService` provider (the FORBIDDEN direct dependency) and a `ChatRuntimeRouter` provider
|
|
* (fronting a real EmbeddedChatRuntime over a SEPARATE scoped AgentService) are registered. Production
|
|
* resolves whichever its constructor declares:
|
|
* - RED today: the controller/gateway `@Inject(AgentService)` → the direct fake is consulted, the
|
|
* router (and its embedded fake) is never reached.
|
|
* - GREEN later: the controller/gateway inject `ChatRuntimeRouter` → the direct fake is never
|
|
* touched (stays at zero) and scope is observed inside the embedded fake behind the router.
|
|
* The SAME test body reds today and greens later; a method-for-method AgentService shim on the router
|
|
* records a FORBIDDEN op (and never the frozen legacy op) in the observed runtime call set, and
|
|
* restoring the direct injection cannot satisfy the "direct fake at zero" / "embedded fake observed
|
|
* scope" / "frozen op invoked on the router" anchors. The router is wrapped by {@link
|
|
* makeRecordingRouter} so those anchors are runtime invocation evidence, not source substrings.
|
|
*/
|
|
function buildRestModule(
|
|
directAgentService: ScopedAgentService,
|
|
embeddedAgentService: ScopedAgentService,
|
|
routerCalls: string[],
|
|
): Promise<TestingModule> {
|
|
return (
|
|
Test.createTestingModule({
|
|
controllers: [ChatController],
|
|
providers: [
|
|
{ provide: AgentService, useValue: directAgentService },
|
|
{
|
|
provide: ChatRuntimeRouter,
|
|
useFactory: () =>
|
|
makeRecordingRouter(legacyRouterFronting(embeddedAgentService), routerCalls),
|
|
},
|
|
],
|
|
})
|
|
// ChatController's @UseGuards(AuthGuard) is resolved during instance loading; AuthGuard injects
|
|
// AUTH, an HTTP-only concern never exercised by a direct handler call. Stub it so the graph
|
|
// resolves and the test reds on BEHAVIOUR, not on a DI collection error.
|
|
.overrideGuard(AuthGuard)
|
|
.useValue({ canActivate: () => true })
|
|
.compile()
|
|
);
|
|
}
|
|
|
|
function buildGatewayModule(
|
|
directAgentService: ScopedAgentService,
|
|
embeddedAgentService: ScopedAgentService,
|
|
routerCalls: string[],
|
|
): Promise<TestingModule> {
|
|
const brain = {
|
|
conversations: {
|
|
// The sender OWNS this durable conversation, so the browser-send admission gate lets the turn
|
|
// reach the router seam. Foreignness is asserted downstream at the in-memory agent session
|
|
// (getSession({USER_B}) -> undefined), not at durable admission — the admission-rejection
|
|
// property has its own dedicated coverage.
|
|
findById: vi.fn().mockResolvedValue({ id: CONVERSATION_ID, userId: USER_B.id }),
|
|
create: vi.fn().mockResolvedValue(undefined),
|
|
update: vi.fn().mockResolvedValue(undefined),
|
|
findMessages: vi.fn().mockResolvedValue([]),
|
|
addMessage: vi.fn().mockResolvedValue({ id: 'persisted-turn' }),
|
|
},
|
|
};
|
|
return Test.createTestingModule({
|
|
providers: [
|
|
ChatGateway,
|
|
{ provide: AgentService, useValue: directAgentService },
|
|
{ provide: AUTH, useValue: { api: { getSession: vi.fn().mockResolvedValue(null) } } },
|
|
{ provide: BRAIN, useValue: brain },
|
|
{ provide: CommandRegistryService, useValue: { getManifest: vi.fn().mockReturnValue([]) } },
|
|
{ provide: CommandExecutorService, useValue: { execute: vi.fn() } },
|
|
{
|
|
provide: RoutingEngineService,
|
|
useValue: {
|
|
resolve: vi.fn().mockResolvedValue({ provider: 'test', model: 'test-model' }),
|
|
},
|
|
},
|
|
{
|
|
provide: ChatRuntimeRouter,
|
|
useFactory: () =>
|
|
makeRecordingRouter(legacyRouterFronting(embeddedAgentService), routerCalls),
|
|
},
|
|
],
|
|
}).compile();
|
|
}
|
|
|
|
describe('TESS-M1-SEC-002 AgentService ownership boundary', () => {
|
|
it('requires explicit owner+tenant scope on protected session operations', () => {
|
|
const source = readFileSync(resolve('src/agent/agent.service.ts'), 'utf8');
|
|
|
|
expect(source).toContain('getSession(sessionId: string, scope: ActorTenantScope)');
|
|
expect(source).toContain('listSessions(scope: ActorTenantScope)');
|
|
expect(source).toContain('getSessionInfo(sessionId: string, scope: ActorTenantScope)');
|
|
expect(source).toContain(
|
|
'addChannel(sessionId: string, channel: string, scope: ActorTenantScope)',
|
|
);
|
|
expect(source).toContain(
|
|
'removeChannel(sessionId: string, channel: string, scope: ActorTenantScope)',
|
|
);
|
|
expect(source).toContain(
|
|
'async prompt(sessionId: string, message: string, scope: ActorTenantScope)',
|
|
);
|
|
expect(source).toContain('scope: ActorTenantScope,');
|
|
expect(source).toContain('async destroySession(sessionId: string, scope: ActorTenantScope)');
|
|
expect(source).not.toContain('scope?: ActorTenantScope');
|
|
});
|
|
});
|
|
|
|
describe('TESS-M1-SEC-002 REST session ownership and tenant binding', () => {
|
|
it('lists only sessions owned by the authenticated owner+tenant scope', () => {
|
|
const agentService = makeScopedAgentService();
|
|
const controller = new SessionsController(agentService as never);
|
|
|
|
expect(controller.list(USER_B)).toEqual({ sessions: [], total: 0 });
|
|
expect(agentService.listSessions).toHaveBeenCalledWith({
|
|
userId: USER_B.id,
|
|
tenantId: USER_B.tenantId,
|
|
});
|
|
});
|
|
|
|
it('does not reveal another owner/tenant session by guessed id', () => {
|
|
const agentService = makeScopedAgentService();
|
|
const controller = new SessionsController(agentService as never);
|
|
|
|
expect(() => controller.findOne(CONVERSATION_ID, USER_B)).toThrow(NotFoundException);
|
|
expect(agentService.getSessionInfo).toHaveBeenCalledWith(CONVERSATION_ID, {
|
|
userId: USER_B.id,
|
|
tenantId: USER_B.tenantId,
|
|
});
|
|
});
|
|
|
|
it('does not terminate another owner/tenant session by guessed id', async () => {
|
|
const agentService = makeScopedAgentService();
|
|
const controller = new SessionsController(agentService as never);
|
|
|
|
await expect(controller.destroy(CONVERSATION_ID, USER_B)).rejects.toBeInstanceOf(
|
|
NotFoundException,
|
|
);
|
|
expect(agentService.destroySession).not.toHaveBeenCalled();
|
|
});
|
|
});
|
|
|
|
describe('TESS-M1-SEC-002 REST chat send ownership and tenant binding (router-delegated legacy runtime)', () => {
|
|
// TESS test A — REST /api/chat send. The genuine RED is the router-delegation redesign, not a slot
|
|
// swap: the forbidden directly-injected AgentService must go UNtouched while the server-derived
|
|
// scope is observed inside the real ChatRuntimeRouter → EmbeddedChatRuntime → AgentService path.
|
|
it('routes a REST send through completeLegacyRestTurn and never the directly-injected AgentService', async () => {
|
|
const directAgentService = makeScopedAgentService(); // FORBIDDEN direct dependency
|
|
const embeddedAgentService = makeScopedAgentService(); // reached ONLY via router → embedded delegation
|
|
const routerCalls: string[] = []; // runtime call set observed AT the controller → router seam
|
|
const moduleRef = await buildRestModule(directAgentService, embeddedAgentService, routerCalls);
|
|
try {
|
|
const controller = moduleRef.get(ChatController, { strict: false });
|
|
|
|
// Foreign ownership is denied (never resolves) — a control that holds today AND at GREEN.
|
|
await expect(
|
|
controller.chat({ conversationId: CONVERSATION_ID, content: 'take over' }, USER_B),
|
|
).rejects.toBeDefined();
|
|
|
|
// Soft anchors so EVERY anchor is evaluated under each mutation, not just the first to fail.
|
|
|
|
// RUNTIME anchor A1 — delegation: the controller must INVOKE the frozen legacy op on the router.
|
|
// Only an actual call enters routerCalls; a dead method/comment/string cannot. RED today (the
|
|
// controller @Inject(AgentService) and never calls the router). GREEN once it drives the op.
|
|
expect
|
|
.soft(routerCalls, 'controller must invoke completeLegacyRestTurn on the router')
|
|
.toContain('completeLegacyRestTurn');
|
|
// RUNTIME anchor A2 — nondelegation: the controller must not drive any AgentService-shaped op on
|
|
// the router. An AgentService-shaped router shim records one of these → RED, defeating the shim
|
|
// on invocation evidence (not source text). A dead named method added alongside the shim does not
|
|
// help: it is never invoked, so it never enters routerCalls while a forbidden op still does.
|
|
for (const op of FORBIDDEN_AGENT_OPS) {
|
|
expect
|
|
.soft(routerCalls, `router seam must not invoke AgentService.${op}`)
|
|
.not.toContain(op);
|
|
}
|
|
// RUNTIME anchor A3 — the forbidden directly-injected AgentService stays at zero (fails today;
|
|
// restoring the direct injection keeps it failing).
|
|
expect.soft(directAgentService.getSession).not.toHaveBeenCalled();
|
|
// RUNTIME anchor A4 — server-derived scope observed INSIDE the separate embedded fake behind the
|
|
// router (fails today; the router path is never taken).
|
|
expect.soft(embeddedAgentService.getSession).toHaveBeenCalledWith(CONVERSATION_ID, {
|
|
userId: USER_B.id,
|
|
tenantId: USER_B.tenantId,
|
|
});
|
|
|
|
// Zero foreign mutation on either path (holds today and at GREEN).
|
|
expect.soft(directAgentService.prompt).not.toHaveBeenCalled();
|
|
expect.soft(embeddedAgentService.prompt).not.toHaveBeenCalled();
|
|
|
|
// Defense-in-depth (NOT load-bearing; the runtime anchors above carry the anti-mask): the
|
|
// controller no longer declares the direct embedded AgentService dependency. A negative source
|
|
// check cannot be satisfied by dead text — it only fails when the injection is present.
|
|
const controllerSource = readFileSync(resolve('src/chat/chat.controller.ts'), 'utf8');
|
|
expect.soft(controllerSource).not.toContain('@Inject(AgentService)');
|
|
} finally {
|
|
await moduleRef.close();
|
|
}
|
|
});
|
|
});
|
|
|
|
describe('TESS-M1-SEC-002 WebSocket session ownership and tenant binding (router-delegated legacy runtime)', () => {
|
|
function makeSocket() {
|
|
return {
|
|
id: 'socket-b',
|
|
connected: true,
|
|
data: { user: USER_B, session: { id: 'auth-session-b', userId: USER_B.id } },
|
|
emit: vi.fn(),
|
|
disconnect: vi.fn(),
|
|
};
|
|
}
|
|
|
|
// TESS test B — WebSocket send/attach.
|
|
it('routes a WebSocket send through prepareLegacySocketTurn and never the directly-injected AgentService', async () => {
|
|
const directAgentService = makeScopedAgentService();
|
|
const embeddedAgentService = makeScopedAgentService();
|
|
const routerCalls: string[] = [];
|
|
const moduleRef = await buildGatewayModule(
|
|
directAgentService,
|
|
embeddedAgentService,
|
|
routerCalls,
|
|
);
|
|
try {
|
|
const gateway = moduleRef.get(ChatGateway, { strict: false });
|
|
const socket = makeSocket();
|
|
|
|
await Promise.resolve(
|
|
gateway.handleMessage(socket as never, {
|
|
conversationId: CONVERSATION_ID,
|
|
content: 'attach to foreign session',
|
|
}),
|
|
).catch(() => undefined);
|
|
|
|
// RUNTIME anchor B1 — delegation: the gateway must invoke the frozen socket op on the router.
|
|
expect
|
|
.soft(routerCalls, 'gateway must invoke prepareLegacySocketTurn on the router')
|
|
.toContain('prepareLegacySocketTurn');
|
|
// RUNTIME anchor B2 — nondelegation: no AgentService-shaped op on the router (defeats the shim).
|
|
for (const op of FORBIDDEN_AGENT_OPS) {
|
|
expect
|
|
.soft(routerCalls, `router seam must not invoke AgentService.${op}`)
|
|
.not.toContain(op);
|
|
}
|
|
// RED anchor B3 — forbidden direct AgentService untouched (fails today, gateway injects it).
|
|
expect.soft(directAgentService.getSession).not.toHaveBeenCalled();
|
|
// RED anchor B4 — scope observed inside router → embedded delegation (fails today, never reached).
|
|
expect.soft(embeddedAgentService.getSession).toHaveBeenCalledWith(CONVERSATION_ID, {
|
|
userId: USER_B.id,
|
|
tenantId: USER_B.tenantId,
|
|
});
|
|
// Foreign session gets zero lease/listener/channel/prompt on EITHER path (holds today and GREEN).
|
|
expect.soft(directAgentService.onEvent).not.toHaveBeenCalled();
|
|
expect.soft(directAgentService.addChannel).not.toHaveBeenCalled();
|
|
expect.soft(directAgentService.prompt).not.toHaveBeenCalled();
|
|
expect.soft(embeddedAgentService.onEvent).not.toHaveBeenCalled();
|
|
expect.soft(embeddedAgentService.addChannel).not.toHaveBeenCalled();
|
|
expect.soft(embeddedAgentService.prompt).not.toHaveBeenCalled();
|
|
expect
|
|
.soft(socket.emit)
|
|
.toHaveBeenCalledWith(
|
|
'error',
|
|
expect.objectContaining({ conversationId: CONVERSATION_ID }),
|
|
);
|
|
|
|
// Defense-in-depth (NOT load-bearing): gateway no longer declares the direct dependency.
|
|
const gatewaySource = readFileSync(resolve('src/chat/chat.gateway.ts'), 'utf8');
|
|
expect.soft(gatewaySource).not.toContain('@Inject(AgentService)');
|
|
} finally {
|
|
await moduleRef.close();
|
|
}
|
|
});
|
|
|
|
// TESS test C — WebSocket set:thinking.
|
|
it('routes set:thinking through setLegacyThinking and never the directly-injected AgentService', async () => {
|
|
const directAgentService = makeScopedAgentService();
|
|
const embeddedAgentService = makeScopedAgentService();
|
|
const routerCalls: string[] = [];
|
|
const moduleRef = await buildGatewayModule(
|
|
directAgentService,
|
|
embeddedAgentService,
|
|
routerCalls,
|
|
);
|
|
try {
|
|
const gateway = moduleRef.get(ChatGateway, { strict: false });
|
|
const socket = makeSocket();
|
|
|
|
await Promise.resolve(
|
|
gateway.handleSetThinking(socket as never, {
|
|
conversationId: CONVERSATION_ID,
|
|
level: 'high',
|
|
}),
|
|
).catch(() => undefined);
|
|
|
|
// RUNTIME anchor C1 — delegation: the gateway must invoke the frozen thinking op on the router.
|
|
expect
|
|
.soft(routerCalls, 'gateway must invoke setLegacyThinking on the router')
|
|
.toContain('setLegacyThinking');
|
|
// RUNTIME anchor C2 — nondelegation: no AgentService-shaped op on the router (defeats the shim).
|
|
for (const op of FORBIDDEN_AGENT_OPS) {
|
|
expect
|
|
.soft(routerCalls, `router seam must not invoke AgentService.${op}`)
|
|
.not.toContain(op);
|
|
}
|
|
expect.soft(directAgentService.getSession).not.toHaveBeenCalled();
|
|
expect.soft(embeddedAgentService.getSession).toHaveBeenCalledWith(CONVERSATION_ID, {
|
|
userId: USER_B.id,
|
|
tenantId: USER_B.tenantId,
|
|
});
|
|
expect
|
|
.soft(socket.emit)
|
|
.toHaveBeenCalledWith(
|
|
'error',
|
|
expect.objectContaining({ conversationId: CONVERSATION_ID }),
|
|
);
|
|
} finally {
|
|
await moduleRef.close();
|
|
}
|
|
});
|
|
|
|
// TESS test D — WebSocket abort.
|
|
it('routes abort through abortLegacyTurn and never the directly-injected AgentService', async () => {
|
|
const directAgentService = makeScopedAgentService();
|
|
const embeddedAgentService = makeScopedAgentService();
|
|
const routerCalls: string[] = [];
|
|
const moduleRef = await buildGatewayModule(
|
|
directAgentService,
|
|
embeddedAgentService,
|
|
routerCalls,
|
|
);
|
|
try {
|
|
const gateway = moduleRef.get(ChatGateway, { strict: false });
|
|
const socket = makeSocket();
|
|
|
|
await Promise.resolve(
|
|
gateway.handleAbort(socket as never, { conversationId: CONVERSATION_ID }),
|
|
).catch(() => undefined);
|
|
|
|
// RUNTIME anchor D1 — delegation: the gateway must invoke the frozen abort op on the router.
|
|
expect
|
|
.soft(routerCalls, 'gateway must invoke abortLegacyTurn on the router')
|
|
.toContain('abortLegacyTurn');
|
|
// RUNTIME anchor D2 — nondelegation: no AgentService-shaped op on the router (defeats the shim).
|
|
for (const op of FORBIDDEN_AGENT_OPS) {
|
|
expect
|
|
.soft(routerCalls, `router seam must not invoke AgentService.${op}`)
|
|
.not.toContain(op);
|
|
}
|
|
expect.soft(directAgentService.getSession).not.toHaveBeenCalled();
|
|
expect.soft(embeddedAgentService.getSession).toHaveBeenCalledWith(CONVERSATION_ID, {
|
|
userId: USER_B.id,
|
|
tenantId: USER_B.tenantId,
|
|
});
|
|
expect
|
|
.soft(socket.emit)
|
|
.toHaveBeenCalledWith(
|
|
'error',
|
|
expect.objectContaining({ conversationId: CONVERSATION_ID }),
|
|
);
|
|
} finally {
|
|
await moduleRef.close();
|
|
}
|
|
});
|
|
|
|
// TESS test E (genuine, unchanged) — pi-rpc browser-legacy refusal.
|
|
it('rejects a browser legacy raw message in pi-rpc mode with a fixed typed unsupported and executes nothing', async () => {
|
|
// pi-rpc: the harness runtime is live. The browser legacy `message` path is unsupported and
|
|
// must be refused with a fixed typed code, touching neither the embedded AgentService nor the
|
|
// harness conversation service.
|
|
const agentService = makeScopedAgentService();
|
|
const embedded = new EmbeddedChatRuntime(agentService as never);
|
|
const harnessConversation = {
|
|
attach: vi.fn(),
|
|
detach: vi.fn(),
|
|
send: vi.fn(),
|
|
subscribeFrom: vi.fn(),
|
|
};
|
|
const harness = new HarnessChatRuntime(harnessConversation as never);
|
|
const router = new ChatRuntimeRouter(
|
|
registryWith(['pi']),
|
|
boundConversationService,
|
|
embedded,
|
|
harness,
|
|
'pi-rpc',
|
|
);
|
|
router.onModuleInit();
|
|
|
|
const brain = {
|
|
conversations: {
|
|
findById: vi.fn().mockResolvedValue(undefined),
|
|
create: vi.fn().mockResolvedValue(undefined),
|
|
update: vi.fn().mockResolvedValue(undefined),
|
|
findMessages: vi.fn().mockResolvedValue([]),
|
|
addMessage: vi.fn().mockResolvedValue(undefined),
|
|
},
|
|
};
|
|
const gateway = new ChatGateway(
|
|
router as never,
|
|
{} as never,
|
|
brain as never,
|
|
{ getManifest: vi.fn().mockReturnValue([]) } as never,
|
|
{ execute: vi.fn() } as never,
|
|
{ resolve: vi.fn() } as never,
|
|
);
|
|
const socket = {
|
|
id: 'socket-b',
|
|
connected: true,
|
|
data: { user: USER_B, session: { id: 'auth-session-b', userId: USER_B.id } },
|
|
emit: vi.fn(),
|
|
disconnect: vi.fn(),
|
|
};
|
|
|
|
await Promise.resolve(
|
|
gateway.handleMessage(socket as never, {
|
|
conversationId: CONVERSATION_ID,
|
|
content: 'route me',
|
|
}),
|
|
).catch(() => undefined);
|
|
|
|
expect(socket.emit).toHaveBeenCalledWith(
|
|
'error',
|
|
expect.objectContaining({ code: 'runtime_unsupported' }),
|
|
);
|
|
expect(agentService.getSession).not.toHaveBeenCalled();
|
|
expect(agentService.prompt).not.toHaveBeenCalled();
|
|
expect(harnessConversation.attach).not.toHaveBeenCalled();
|
|
expect(harnessConversation.send).not.toHaveBeenCalled();
|
|
});
|
|
});
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// Task-5 AMEND — embedded runtime lease lifecycle (G1) + ownership collapse (G5).
|
|
// These drive the real EmbeddedChatRuntime directly over a shape-complete AgentService
|
|
// fake (every touched method exists, so a RED can only come from behavior, never a
|
|
// `getSession is not a function` TypeError). Ownership context is minted through the
|
|
// real `ownConversation` factory — the only sanctioned way to reach a port op.
|
|
// ---------------------------------------------------------------------------
|
|
|
|
const EMBEDDED_SCOPE = { userId: USER_A.id, tenantId: USER_A.tenantId };
|
|
const CONVERSATION_UNAVAILABLE_RESULT = {
|
|
ok: false,
|
|
code: 'conversation_unavailable',
|
|
retryable: false,
|
|
} as const;
|
|
|
|
/** A stream sink; `channelId` is server-derived, `onEvent` records nothing here. */
|
|
function makeStream(): LegacyRuntimeStream {
|
|
return { channelId: 'websocket:test-1', onEvent: vi.fn() };
|
|
}
|
|
|
|
/**
|
|
* getSession → undefined (session missing), createSession → rejects with `err`. Exercises the
|
|
* `resolveOrCreate` collapse branch. `prompt` exists so its ABSENCE from the call record proves
|
|
* the turn short-circuited before any dispatch.
|
|
*/
|
|
function makeCollapsingAgentService(err: Error) {
|
|
return {
|
|
getSession: vi.fn(() => undefined),
|
|
createSession: vi.fn().mockRejectedValue(err),
|
|
onEvent: vi.fn(() => vi.fn()),
|
|
addChannel: vi.fn(),
|
|
removeChannel: vi.fn(),
|
|
prompt: vi.fn().mockResolvedValue(undefined),
|
|
recordTokenUsage: vi.fn(),
|
|
};
|
|
}
|
|
|
|
/** getSession → a live owned session, so `resolveOrCreate` succeeds and a lease is built. */
|
|
function makeLeaseAgentService() {
|
|
const session = makeAgentSession(USER_A);
|
|
const unsubscribe = vi.fn();
|
|
const svc = {
|
|
getSession: vi.fn(() => session),
|
|
createSession: vi.fn(),
|
|
onEvent: vi.fn(() => unsubscribe),
|
|
addChannel: vi.fn(),
|
|
removeChannel: vi.fn(),
|
|
prompt: vi.fn().mockResolvedValue(undefined),
|
|
recordTokenUsage: vi.fn(),
|
|
};
|
|
return { svc, unsubscribe, session };
|
|
}
|
|
|
|
/**
|
|
* getSession → a live owned session (REST resolveOrCreate succeeds), onEvent returns a `detach`
|
|
* spy, and `prompt` REJECTS with a non-timeout error. Drives the REST-turn catch path so the single
|
|
* idempotent teardown must clear the 120s timeout and detach the listener exactly once.
|
|
*/
|
|
function makeRejectingPromptAgentService() {
|
|
const session = makeAgentSession(USER_A);
|
|
const detach = vi.fn();
|
|
const svc = {
|
|
getSession: vi.fn(() => session),
|
|
createSession: vi.fn(),
|
|
onEvent: vi.fn(() => detach),
|
|
addChannel: vi.fn(),
|
|
removeChannel: vi.fn(),
|
|
prompt: vi.fn().mockRejectedValue(new Error('agent backend exploded')),
|
|
recordTokenUsage: vi.fn(),
|
|
};
|
|
return { svc, detach };
|
|
}
|
|
|
|
describe('TESS Task-5 embedded ownership collapse (missing and foreign are indistinguishable, never throw)', () => {
|
|
const ctx = ownConversation(CONVERSATION_ID, EMBEDDED_SCOPE);
|
|
|
|
it('collapses a foreign (Forbidden) create to conversation_unavailable and never throws', async () => {
|
|
const svc = makeCollapsingAgentService(new ForbiddenException('foreign owner'));
|
|
const runtime = new EmbeddedChatRuntime(svc as never);
|
|
|
|
const result = await runtime.completeLegacyRestTurn(ctx, { content: 'take over' });
|
|
|
|
expect(result).toEqual(CONVERSATION_UNAVAILABLE_RESULT);
|
|
expect(svc.prompt).not.toHaveBeenCalled();
|
|
});
|
|
|
|
it('collapses a missing (NotFound) create to conversation_unavailable and never throws', async () => {
|
|
const svc = makeCollapsingAgentService(new NotFoundException('no such conversation'));
|
|
const runtime = new EmbeddedChatRuntime(svc as never);
|
|
|
|
const result = await runtime.completeLegacyRestTurn(ctx, { content: 'hello' });
|
|
|
|
expect(result).toEqual(CONVERSATION_UNAVAILABLE_RESULT);
|
|
expect(svc.prompt).not.toHaveBeenCalled();
|
|
});
|
|
|
|
it('returns the IDENTICAL collapse for foreign and missing so neither can be distinguished', async () => {
|
|
const foreign = new EmbeddedChatRuntime(
|
|
makeCollapsingAgentService(new ForbiddenException('foreign owner')) as never,
|
|
);
|
|
const missing = new EmbeddedChatRuntime(
|
|
makeCollapsingAgentService(new NotFoundException('no such conversation')) as never,
|
|
);
|
|
|
|
const foreignResult = await foreign.completeLegacyRestTurn(ctx, { content: 'x' });
|
|
const missingResult = await missing.completeLegacyRestTurn(ctx, { content: 'x' });
|
|
|
|
expect(foreignResult).toEqual(missingResult);
|
|
expect(foreignResult).toEqual(CONVERSATION_UNAVAILABLE_RESULT);
|
|
});
|
|
});
|
|
|
|
describe('TESS Task-5 embedded socket lease lifecycle (one-shot dispatch, idempotent dispose, partial-setup rollback)', () => {
|
|
const ctx = ownConversation(CONVERSATION_ID, EMBEDDED_SCOPE);
|
|
|
|
it('dispatches the turn exactly once; a second dispatch is a no-op turn_already_dispatched', async () => {
|
|
const { svc } = makeLeaseAgentService();
|
|
const runtime = new EmbeddedChatRuntime(svc as never);
|
|
|
|
const prepared = await runtime.prepareLegacySocketTurn(ctx, { content: 'first' }, makeStream());
|
|
expect(prepared.ok).toBe(true);
|
|
if (!prepared.ok) throw new Error('prepareLegacySocketTurn should succeed');
|
|
const lease = prepared.value;
|
|
|
|
const first = await lease.dispatch();
|
|
expect(first).toEqual({ ok: true, value: undefined });
|
|
expect(svc.prompt).toHaveBeenCalledTimes(1);
|
|
|
|
const second = await lease.dispatch();
|
|
expect(second).toEqual({ ok: false, code: 'turn_already_dispatched', retryable: false });
|
|
// Zero additional effect — the second dispatch must not prompt again.
|
|
expect(svc.prompt).toHaveBeenCalledTimes(1);
|
|
});
|
|
|
|
it('disposes once; a second dispose is a silent no-op that never re-detaches or destroys the session', async () => {
|
|
const { svc, unsubscribe, session } = makeLeaseAgentService();
|
|
const runtime = new EmbeddedChatRuntime(svc as never);
|
|
|
|
const prepared = await runtime.prepareLegacySocketTurn(ctx, { content: 'x' }, makeStream());
|
|
expect(prepared.ok).toBe(true);
|
|
if (!prepared.ok) throw new Error('prepareLegacySocketTurn should succeed');
|
|
const lease = prepared.value;
|
|
|
|
await lease.dispose();
|
|
await lease.dispose();
|
|
|
|
// Listener + channel torn down exactly once across two dispose calls.
|
|
expect(unsubscribe).toHaveBeenCalledTimes(1);
|
|
expect(svc.removeChannel).toHaveBeenCalledTimes(1);
|
|
// Disposal never terminates the underlying session or process.
|
|
expect(session.piSession.abort).not.toHaveBeenCalled();
|
|
expect(session.piSession.dispose).not.toHaveBeenCalled();
|
|
});
|
|
|
|
it('rolls back the acquired listener and returns a total safe failure when channel attach fails mid-setup', async () => {
|
|
const { svc, unsubscribe } = makeLeaseAgentService();
|
|
svc.addChannel = vi.fn(() => {
|
|
throw new Error('channel attach failed');
|
|
});
|
|
const runtime = new EmbeddedChatRuntime(svc as never);
|
|
|
|
// Must NOT throw out of the port — a partial setup collapses to a total safe failure.
|
|
const prepared = await runtime.prepareLegacySocketTurn(ctx, { content: 'x' }, makeStream());
|
|
expect(prepared.ok).toBe(false);
|
|
// Exactly what was acquired (the event listener) is rolled back.
|
|
expect(unsubscribe).toHaveBeenCalledTimes(1);
|
|
});
|
|
});
|
|
|
|
describe('TESS Task-5 embedded REST turn teardown (a prompt rejection frees the timer + listener exactly once)', () => {
|
|
const ctx = ownConversation(CONVERSATION_ID, EMBEDDED_SCOPE);
|
|
|
|
it('clears the 120s timeout and detaches the listener exactly once when prompt() rejects, leaving no timer to reject the abandoned done-promise later (Task 5 finding 6)', async () => {
|
|
const { svc, detach } = makeRejectingPromptAgentService();
|
|
const runtime = new EmbeddedChatRuntime(svc as never);
|
|
|
|
// A rejected `done` promise firing after completeLegacyRestTurn has already returned would
|
|
// surface as an unhandledRejection — the leak this test fences. Capture any that escape.
|
|
const unhandled: unknown[] = [];
|
|
const onUnhandled = (reason: unknown): void => {
|
|
unhandled.push(reason);
|
|
};
|
|
process.on('unhandledRejection', onUnhandled);
|
|
vi.useFakeTimers();
|
|
try {
|
|
const result = await runtime.completeLegacyRestTurn(ctx, {
|
|
content: 'trigger a backend failure',
|
|
});
|
|
|
|
// The rejection collapses to a total safe failure (not a timeout) — never throws out of the port.
|
|
expect(result).toEqual({ ok: false, code: 'operation_failed', retryable: false });
|
|
// The single idempotent dispose ran in the catch: listener detached exactly once.
|
|
expect(detach).toHaveBeenCalledTimes(1);
|
|
|
|
// dispose() cleared the REST timeout, so advancing far past it (120s) fires nothing: no second
|
|
// detach, and — the actual leak — no live timer left to reject the now-abandoned `done` promise.
|
|
vi.advanceTimersByTime(600_000);
|
|
expect(detach).toHaveBeenCalledTimes(1);
|
|
} finally {
|
|
vi.useRealTimers();
|
|
}
|
|
// Let any scheduled rejection surface on a real macrotask, then confirm none did.
|
|
await new Promise((resolve) => setTimeout(resolve, 0));
|
|
process.off('unhandledRejection', onUnhandled);
|
|
expect(unhandled).toHaveLength(0);
|
|
});
|
|
|
|
it('bounds a hung prompt: when prompt() never settles and no agent_end arrives, the 120s timeout ends the turn with a timeout result and exactly one teardown, no unhandledRejection (Task 5 finding 6 — pending-prompt timeout)', async () => {
|
|
const session = makeAgentSession(USER_A);
|
|
const detach = vi.fn();
|
|
const svc = {
|
|
getSession: vi.fn(() => session),
|
|
createSession: vi.fn(),
|
|
onEvent: vi.fn(() => detach),
|
|
addChannel: vi.fn(),
|
|
removeChannel: vi.fn(),
|
|
// The prompt never resolves or rejects — a hung agent backend. Under the pre-fix sequential
|
|
// `await prompt()` the timer could never even be observed, so the turn hung forever.
|
|
prompt: vi.fn(() => new Promise<void>(() => undefined)),
|
|
recordTokenUsage: vi.fn(),
|
|
};
|
|
const runtime = new EmbeddedChatRuntime(svc as never);
|
|
|
|
const unhandled: unknown[] = [];
|
|
const onUnhandled = (reason: unknown): void => {
|
|
unhandled.push(reason);
|
|
};
|
|
process.on('unhandledRejection', onUnhandled);
|
|
vi.useFakeTimers();
|
|
try {
|
|
const resultPromise = runtime.completeLegacyRestTurn(ctx, {
|
|
content: 'a prompt that never returns',
|
|
});
|
|
// No agent_end, prompt still pending: only the 120s timeout can end the turn. Promise.all
|
|
// installed a handler on `done` synchronously, so the timer bounds the turn while prompt hangs.
|
|
await vi.advanceTimersByTimeAsync(200_000);
|
|
const result = await resultPromise;
|
|
|
|
expect(result).toEqual({ ok: false, code: 'timeout', retryable: true });
|
|
// The single idempotent dispose ran on the timeout path: listener detached exactly once.
|
|
expect(detach).toHaveBeenCalledTimes(1);
|
|
// Advancing far past the deadline fires nothing more: dispose cleared the timer.
|
|
vi.advanceTimersByTime(600_000);
|
|
expect(detach).toHaveBeenCalledTimes(1);
|
|
} finally {
|
|
vi.useRealTimers();
|
|
}
|
|
await new Promise((resolve) => setTimeout(resolve, 0));
|
|
process.off('unhandledRejection', onUnhandled);
|
|
expect(unhandled).toHaveLength(0);
|
|
});
|
|
|
|
it('when the 120s timeout fires while prompt() is still pending, returns timeout with one teardown, and a later prompt rejection surfaces no unhandledRejection (Task 5 finding 6 — timeout/prompt race)', async () => {
|
|
const session = makeAgentSession(USER_A);
|
|
const detach = vi.fn();
|
|
let rejectPrompt: (reason: unknown) => void = () => undefined;
|
|
const prompting = new Promise<void>((_resolve, reject) => {
|
|
rejectPrompt = reject;
|
|
});
|
|
const svc = {
|
|
getSession: vi.fn(() => session),
|
|
createSession: vi.fn(),
|
|
onEvent: vi.fn(() => detach),
|
|
addChannel: vi.fn(),
|
|
removeChannel: vi.fn(),
|
|
prompt: vi.fn(() => prompting),
|
|
recordTokenUsage: vi.fn(),
|
|
};
|
|
const runtime = new EmbeddedChatRuntime(svc as never);
|
|
|
|
const unhandled: unknown[] = [];
|
|
const onUnhandled = (reason: unknown): void => {
|
|
unhandled.push(reason);
|
|
};
|
|
process.on('unhandledRejection', onUnhandled);
|
|
vi.useFakeTimers();
|
|
try {
|
|
const resultPromise = runtime.completeLegacyRestTurn(ctx, {
|
|
content: 'prompt settles after the deadline',
|
|
});
|
|
// The timeout wins the race while prompt is still pending.
|
|
await vi.advanceTimersByTimeAsync(200_000);
|
|
const result = await resultPromise;
|
|
|
|
expect(result).toEqual({ ok: false, code: 'timeout', retryable: true });
|
|
expect(detach).toHaveBeenCalledTimes(1);
|
|
|
|
// The prompt now rejects LATE — after the turn already returned its timeout result. Because
|
|
// Promise.all installed a rejection handler on `prompting` synchronously (the fix), this late
|
|
// rejection is already observed and must not escape as an unhandledRejection.
|
|
rejectPrompt(new Error('late backend failure'));
|
|
} finally {
|
|
vi.useRealTimers();
|
|
}
|
|
await new Promise((resolve) => setTimeout(resolve, 0));
|
|
process.off('unhandledRejection', onUnhandled);
|
|
expect(unhandled).toHaveLength(0);
|
|
});
|
|
});
|