From 146d19f3e59d5a8c098d74d138d99ef459a17d0d Mon Sep 17 00:00:00 2001 From: Justin Middler Date: Mon, 27 Jul 2026 17:39:06 +1000 Subject: [PATCH] Preserve Slack harness context --- .../src/runtime.integration.test.ts | 180 +++++++++++++++++- apps/agent-gateway/src/runtime.ts | 33 ++-- 2 files changed, 199 insertions(+), 14 deletions(-) diff --git a/apps/agent-gateway/src/runtime.integration.test.ts b/apps/agent-gateway/src/runtime.integration.test.ts index 8a4766b..d1f1d84 100644 --- a/apps/agent-gateway/src/runtime.integration.test.ts +++ b/apps/agent-gateway/src/runtime.integration.test.ts @@ -1,22 +1,39 @@ import { createHash } from "node:crypto"; +import { createServer } from "node:http"; import { afterAll, beforeAll, describe, expect, it } from "vitest"; import { closeDatabase, database, newId, schema } from "@muster/database"; import { AgentStructuredOutputSchemas, HuntResultSchema, } from "@muster/contracts"; +import { encryptConnectorAuth } from "@muster/integrations"; import { and, desc, eq } from "drizzle-orm"; import { z } from "zod"; import { bindHuntResultToAuthoritativeCase, codexOutputSchemaFor, DurableAgentRuntime, + parsePersistedRequest, } from "./runtime.ts"; const integration = process.env.MUSTER_INTEGRATION_TESTS === "true"; const describeIntegration = integration ? describe.sequential : describe.skip; describe("Codex structured output schema", () => { + it("preserves Slack harness mode for live connector context", () => { + expect( + parsePersistedRequest({ + kind: "direct_message", + humanRequest: "Which Tawny hosts need attention?", + harness: { mode: "slack" }, + }), + ).toMatchObject({ + kind: "direct_message", + humanRequest: "Which Tawny hosts need attention?", + harness: { mode: "slack" }, + }); + }); + it("removes unsupported URI formats while preserving authoritative validation", () => { const generated = z.toJSONSchema(AgentStructuredOutputSchemas.HuntResult, { target: "draft-2020-12", @@ -169,7 +186,10 @@ describeIntegration("durable agent runtime", () => { throw new Error(`Run ${runId} did not reach ${status}`); } - async function directMessageSource(suffix: string) { + async function directMessageSource( + suffix: string, + targetAgentId = agentId, + ) { const [room] = await database() .select({ id: schema.rooms.id }) .from(schema.rooms) @@ -178,7 +198,7 @@ describeIntegration("durable agent runtime", () => { and( eq(schema.roomMemberships.organisationId, organisationId), eq(schema.roomMemberships.roomId, schema.rooms.id), - eq(schema.roomMemberships.actorId, agentId), + eq(schema.roomMemberships.actorId, targetAgentId), ), ) .where( @@ -426,6 +446,162 @@ describeIntegration("durable agent runtime", () => { costRuntime.stop(); }); + it("loads live connector evidence for Slack runs", async () => { + const [jessie] = await database() + .select() + .from(schema.agentDefinitions) + .where(eq(schema.agentDefinitions.name, "Jessie")) + .limit(1); + if (!jessie) throw new Error("Bootstrapped Jessie required"); + const source = await directMessageSource("slack-live-context", jessie.id); + const connectorServer = createServer((_request, response) => { + response.writeHead(200, { "content-type": "application/json" }); + response.end( + JSON.stringify([ + { + id: "synthetic-host-20", + hostname: "synthetic-host-20.example.test", + status: "online", + }, + ]), + ); + }); + await new Promise((resolve) => + connectorServer.listen(0, "127.0.0.1", resolve), + ); + const address = connectorServer.address(); + if (!address || typeof address === "string") + throw new Error("Synthetic connector port unavailable"); + const integrationId = newId(); + const templateId = newId(); + const encryptionKey = `synthetic-connector-key-${newId()}`; + const previousEncryptionKey = process.env.CONNECTOR_ENCRYPTION_KEY; + let runtime: DurableAgentRuntime | undefined; + + try { + await database() + .insert(schema.integrationRecords) + .values({ + id: integrationId, + organisationId, + product: "tawny", + instanceId: `runtime-slack-${integrationId}`, + displayName: "Synthetic live Tawny", + status: "configured", + mock: false, + configuration: { + product: "tawny", + instanceId: `runtime-slack-${integrationId}`, + displayName: "Synthetic live Tawny", + baseUrl: `http://127.0.0.1:${address.port}`, + allowedHosts: ["127.0.0.1"], + allowPrivateNetwork: true, + testMode: true, + authType: "none", + limits: { + timeoutMs: 500, + maxResponseBytes: 4_096, + maxRecords: 10, + maxPages: 1, + requestsPerMinute: 60, + }, + }, + }); + await database() + .insert(schema.integrationQueryTemplates) + .values({ + id: templateId, + organisationId, + integrationId, + templateKey: "tawny.inventory.list", + version: 1, + definition: { + key: "tawny.inventory.list", + version: 1, + displayName: "Synthetic Tawny inventory", + method: "GET", + pathTemplate: "/api/agents", + requiredCapability: "tawny.telemetry.read", + inputSchema: { + type: "object", + additionalProperties: false, + }, + outputSchema: { + type: "array", + items: { type: "object" }, + }, + }, + createdByActorId: requestedByActorId, + }); + await database() + .insert(schema.integrationConnectorCredentials) + .values({ + organisationId, + integrationId, + encryptedCredential: encryptConnectorAuth( + { type: "none" }, + encryptionKey, + ), + rotatedByActorId: requestedByActorId, + }); + process.env.CONNECTOR_ENCRYPTION_KEY = encryptionKey; + const run = await insertRun("slack-live-context", { + agentId: jessie.id, + roomId: source.roomId, + promptVersion: jessie.systemPromptVersion, + request: { + kind: "direct_message", + sourceMessageId: source.messageId, + humanRequest: "Which Tawny hosts need attention?", + traceId: `integration-slack-${source.messageId}`, + harness: { mode: "slack" }, + }, + }); + runtime = new DurableAgentRuntime({ + executionRuntime: "mock", + codexHome: "/tmp/muster-runtime-integration", + mockDelayMs: 10, + }); + await runtime.dispatch(); + await waitFor(run.id, "completed"); + + const [query] = await database() + .select() + .from(schema.integrationQueryRuns) + .where( + and( + eq(schema.integrationQueryRuns.organisationId, organisationId), + eq(schema.integrationQueryRuns.integrationId, integrationId), + ), + ) + .limit(1); + expect(query).toMatchObject({ + status: "succeeded", + requestedByActorId: jessie.id, + result: [ + { + id: "synthetic-host-20", + hostname: "synthetic-host-20.example.test", + status: "online", + }, + ], + requestMetadata: { + source: "agent-live-context", + agentRunId: run.id, + templateKey: "tawny.inventory.list", + }, + }); + } finally { + runtime?.stop(); + if (previousEncryptionKey === undefined) + delete process.env.CONNECTOR_ENCRYPTION_KEY; + else process.env.CONNECTOR_ENCRYPTION_KEY = previousEncryptionKey; + await new Promise((resolve) => + connectorServer.close(() => resolve()), + ); + } + }); + it("correlates governed hunt evidence without obeying connector prompt injection", async () => { const [jessie] = await database() .select() diff --git a/apps/agent-gateway/src/runtime.ts b/apps/agent-gateway/src/runtime.ts index cf635e8..309356f 100644 --- a/apps/agent-gateway/src/runtime.ts +++ b/apps/agent-gateway/src/runtime.ts @@ -56,9 +56,28 @@ type PersistedRequest = { traceId?: string | undefined; harness?: { mode?: "slack" | "hermes" | "mcp" | "cli" | "http" | undefined; - }; + } | undefined; }; +const PersistedRequestSchema = z.object({ + kind: z.enum(["jessie_hunt", "direct_message"]).optional(), + huntId: z.uuid().optional(), + huntPlan: z.unknown().optional(), + humanRequest: z.string().optional(), + sourceMessageId: z.uuid().optional(), + traceId: z.string().optional(), + harness: z + .object({ + mode: z.enum(["slack", "hermes", "mcp", "cli", "http"]).optional(), + }) + .optional(), +}); + +export function parsePersistedRequest(input: unknown): PersistedRequest { + const parsed = PersistedRequestSchema.safeParse(input); + return parsed.success ? parsed.data : {}; +} + type LiveConnectorEvidence = { queryRunId?: string; source: string; @@ -1704,17 +1723,7 @@ export class DurableAgentRuntime { } private request(run: AgentRunRow): PersistedRequest { - const parsed = z - .object({ - kind: z.enum(["jessie_hunt", "direct_message"]).optional(), - huntId: z.uuid().optional(), - huntPlan: z.unknown().optional(), - humanRequest: z.string().optional(), - sourceMessageId: z.uuid().optional(), - traceId: z.string().optional(), - }) - .safeParse(run.request); - return parsed.success ? parsed.data : {}; + return parsePersistedRequest(run.request); } private job(run: AgentRunRow): AgentInvestigationJob {