Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -21,3 +21,6 @@ __pycache__/
docs/master.html

.context/

adrs/qm-grok-bridge.md
films/
334 changes: 334 additions & 0 deletions docs/grok-bridge.md

Large diffs are not rendered by default.

2 changes: 2 additions & 0 deletions src/api/deps.ts
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,7 @@ import type { WebhookReceiver } from "../webhooks/webhook-receiver.ts";
import type { IdentityService } from "../identity/identity-service.ts";
import type { DeviceFlowCutoverStore } from "../credentials/device-flow-cutover.ts";
import type { FeatureFlagStore } from "../feature-flags.ts";
import type { GrokBridge } from "../grok-bridge/types.ts";
import type {
ConnectorTokenStore,
Keychain,
Expand Down Expand Up @@ -99,6 +100,7 @@ export interface ServerDeps {
credentialUsage?: CredentialUsageSink;
deviceFlowCutover?: DeviceFlowCutoverStore;
featureFlags?: FeatureFlagStore;
grokBridge?: GrokBridge;
egressAudit?: EgressAuditSink;
brokerFetch?: BrokerFetch;
gitHttpFetch?: GitHttpFetch;
Expand Down
204 changes: 204 additions & 0 deletions src/api/routes/grok-bridge.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,204 @@
import { createManualProvisioner } from "../../grok-bridge/provisioner.ts";
import { GrokBridgeError } from "../../grok-bridge/types.ts";
import { personalScope } from "../../types.ts";
import { errMessage } from "../../util/errors.ts";
import { PayloadTooLargeError, readRawBody, sendJson } from "../http.ts";
import { isObj, orgScope } from "./shared.ts";
import type { ApiCtx, BaseCtx, Route } from "./route.ts";

const provisioner = createManualProvisioner();

function actorIdOf(ctx: ApiCtx, body: Record<string, unknown>): string | undefined {
if (ctx.capability?.actorId) return ctx.capability.actorId;
if (ctx.actor?.p) return ctx.actor.p;
return typeof body.actorId === "string" ? body.actorId : undefined;
}

function actorTypeOf(body: Record<string, unknown>): "internal" | "guest" {
return body.actorType === "guest" ? "guest" : "internal";
}

async function requireBridge(ctx: ApiCtx | BaseCtx): Promise<boolean> {
const bridge = ctx.deps.grokBridge;
const flags = ctx.deps.featureFlags;
if (!bridge || !flags) {
sendJson(ctx.res, 404, { error: "not_found" });
return false;
}
const on = await flags.enabled("grok_bridge", orgScope());
if (!on) {
sendJson(ctx.res, 404, { error: "not_found" });
return false;
}
return true;
}

function sendBridgeError(ctx: { res: ApiCtx["res"] }, error: unknown): void {
if (error instanceof GrokBridgeError) {
sendJson(ctx.res, error.status, { error: error.code, message: error.message });
return;
}
sendJson(ctx.res, 500, { error: "internal", message: errMessage(error) });
}

async function createPairing(ctx: ApiCtx): Promise<void> {
if (!(await requireBridge(ctx))) return;
const body = isObj(ctx.body) ? ctx.body : {};
const actorId = actorIdOf(ctx, body);
if (!actorId || typeof body.agentName !== "string" || typeof body.ownerPrincipalId !== "string") {
return sendJson(ctx.res, 400, {
error: "bad_request",
message: "agentName, ownerPrincipalId, and actor are required",
});
}
try {
const pairing = await ctx.deps.grokBridge!.requestPairing({
agentName: body.agentName,
ownerPrincipalId: body.ownerPrincipalId,
actorId,
actorType: actorTypeOf(body),
originScopeId: typeof body.originScopeId === "string" ? body.originScopeId : personalScope(body.ownerPrincipalId),
});
sendJson(ctx.res, 200, {
pairing,
skill: provisioner.skillFor(pairing),
});
} catch (error) {
sendBridgeError(ctx, error);
}
}

async function decidePairing(ctx: ApiCtx): Promise<void> {
if (!(await requireBridge(ctx))) return;
const body = isObj(ctx.body) ? ctx.body : {};
const actorId = actorIdOf(ctx, body);
if (!actorId || (body.decision !== "accept" && body.decision !== "decline")) {
return sendJson(ctx.res, 400, { error: "bad_request", message: "decision must be accept or decline" });
}
try {
sendJson(ctx.res, 200, {
pairing: await ctx.deps.grokBridge!.decidePairing(ctx.params.id!, actorId, body.decision),
});
} catch (error) {
sendBridgeError(ctx, error);
}
}

async function completeInbound(ctx: ApiCtx): Promise<void> {
if (!(await requireBridge(ctx))) return;
const body = isObj(ctx.body) ? ctx.body : {};
const actorId = actorIdOf(ctx, body);
if (!actorId || typeof body.webhookUrl !== "string" || typeof body.webhookKey !== "string") {
return sendJson(ctx.res, 400, { error: "bad_request", message: "webhookUrl and webhookKey are required" });
}
try {
sendJson(ctx.res, 200, {
pairing: await ctx.deps.grokBridge!.completeInbound(ctx.params.id!, actorId, {
webhookUrl: body.webhookUrl,
webhookKey: body.webhookKey,
...(typeof body.grokBotId === "string" ? { grokBotId: body.grokBotId } : {}),
}),
});
} catch (error) {
sendBridgeError(ctx, error);
}
}

async function revokePairing(ctx: ApiCtx): Promise<void> {
if (!(await requireBridge(ctx))) return;
const body = isObj(ctx.body) ? ctx.body : {};
const actorId = actorIdOf(ctx, body);
if (!actorId) return sendJson(ctx.res, 400, { error: "bad_request", message: "actor is required" });
try {
await ctx.deps.grokBridge!.revoke(ctx.params.id!, actorId);
sendJson(ctx.res, 200, { ok: true });
} catch (error) {
sendBridgeError(ctx, error);
}
}

async function createJob(ctx: ApiCtx): Promise<void> {
if (!(await requireBridge(ctx))) return;
const body = isObj(ctx.body) ? ctx.body : {};
const actorId = actorIdOf(ctx, body);
if (
!actorId ||
typeof body.agentName !== "string" ||
typeof body.ownerPrincipalId !== "string" ||
typeof body.originSessionId !== "string" ||
typeof body.instruction !== "string"
) {
return sendJson(ctx.res, 400, {
error: "bad_request",
message: "agentName, ownerPrincipalId, originSessionId, instruction, and actor are required",
});
}
const publicBase = ctx.deps.publicUrl ?? ctx.deps.portalUrl ?? "";
try {
sendJson(ctx.res, 200, {
job: await ctx.deps.grokBridge!.dispatch({
agentName: body.agentName,
ownerPrincipalId: body.ownerPrincipalId,
originSessionId: body.originSessionId,
originActorId: actorId,
actorType: actorTypeOf(body),
instruction: body.instruction,
callbackBaseUrl: publicBase || "http://127.0.0.1",
}),
});
} catch (error) {
sendBridgeError(ctx, error);
}
}

async function getJob(ctx: ApiCtx): Promise<void> {
if (!(await requireBridge(ctx))) return;
const viewer = actorIdOf(ctx, isObj(ctx.body) ? ctx.body : {}) ?? ctx.url.searchParams.get("viewer") ?? undefined;
if (!viewer) return sendJson(ctx.res, 400, { error: "bad_request", message: "viewer is required" });
try {
sendJson(ctx.res, 200, { job: await ctx.deps.grokBridge!.getJob(ctx.params.id!, viewer) });
} catch (error) {
sendBridgeError(ctx, error);
}
}

async function ingestEvent(ctx: BaseCtx): Promise<void> {
const { req, res, deps, params } = ctx;
if (!(await requireBridge(ctx))) return;
let rawBody: string;
try {
rawBody = await readRawBody(req);
} catch (error) {
if (error instanceof PayloadTooLargeError) {
return sendJson(res, 413, { error: "payload_too_large", message: errMessage(error) });
}
return sendJson(res, 400, { error: "bad_request" });
}
let parsed: unknown;
try {
parsed = JSON.parse(rawBody);
} catch {
return sendJson(res, 400, { error: "bad_request", message: "invalid JSON" });
}
const header = req.headers.authorization;
const authorization = Array.isArray(header) ? header[0] : header;
try {
const result = await deps.grokBridge!.ingest(params.id!, parsed, authorization);
sendJson(res, result.duplicate ? 200 : 202, { ok: true, duplicate: result.duplicate });
} catch (error) {
sendBridgeError(ctx, error);
}
}

export const grokBridgeRawRoutes: ReadonlyArray<Route<BaseCtx>> = [
{ method: "POST", path: "/v1/grok-bridge/jobs/:id/events", auth: "public", handle: ingestEvent },
];

export const grokBridgeRoutes: ReadonlyArray<Route<ApiCtx>> = [
{ method: "POST", path: "/v1/grok-bridge/pairings", auth: "either", handle: createPairing },
{ method: "POST", path: "/v1/grok-bridge/pairings/:id/decide", auth: "either", handle: decidePairing },
{ method: "POST", path: "/v1/grok-bridge/pairings/:id/inbound", auth: "either", handle: completeInbound },
{ method: "POST", path: "/v1/grok-bridge/pairings/:id/revoke", auth: "either", handle: revokePairing },
{ method: "POST", path: "/v1/grok-bridge/jobs", auth: "either", handle: createJob },
{ method: "GET", path: "/v1/grok-bridge/jobs/:id", auth: "either", handle: getJob },
];
3 changes: 3 additions & 0 deletions src/api/routes/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@ import { authBrokerRoutes } from "./auth-broker.ts";
import { loopItemRoutes } from "./loop-items.ts";
import { searchRoutes } from "./search.ts";
import { userModelAuthRoutes } from "./user-model-auth.ts";
import { grokBridgeRawRoutes, grokBridgeRoutes } from "./grok-bridge.ts";

export const rawRoutes: ReadonlyArray<Route<BaseCtx>> = [
{ method: "GET", path: "/healthz", auth: "public", handle: ({ res }) => sendJson(res, 200, { ok: true }) },
Expand All @@ -49,6 +50,7 @@ export const rawRoutes: ReadonlyArray<Route<BaseCtx>> = [
...sessionStateRawRoutes,
...loopItemEventsRawRoutes,
...webhookRawRoutes,
...grokBridgeRawRoutes,
];

export const apiRoutes: ReadonlyArray<Route<ApiCtx>> = [
Expand Down Expand Up @@ -80,4 +82,5 @@ export const apiRoutes: ReadonlyArray<Route<ApiCtx>> = [
...egressAuditRoutes,
...authBrokerRoutes,
...userModelAuthRoutes,
...grokBridgeRoutes,
];
2 changes: 1 addition & 1 deletion src/feature-flags.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@ import type { DurableMap } from "./persistence/durable-map.ts";
import { orgId as configOrgId } from "./config.ts";
import { scopeId, type ScopeId } from "./types.ts";

export const FEATURE_NAMES = ["command_scoped_credentials"] as const;
export const FEATURE_NAMES = ["command_scoped_credentials", "grok_bridge"] as const;
export type FeatureName = (typeof FEATURE_NAMES)[number];

export interface FeatureFlagRecord {
Expand Down
29 changes: 29 additions & 0 deletions src/grok-bridge/authz.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,29 @@
import { normalizeAgentName } from "./crypto.ts";
import { GrokBridgeError, type GrokPairing } from "./types.ts";

export function requireOwner(pairing: GrokPairing, ownerId: string): void {
if (pairing.ownerPrincipalId !== ownerId) {
throw new GrokBridgeError("forbidden", 403, "only the pairing owner can do that");
}
}

export function requireInternalActor(actorType: "internal" | "guest", action: string): void {
if (actorType === "guest") {
throw new GrokBridgeError("forbidden", 403, `guests cannot ${action}`);
}
}

export function requireAgentName(raw: string): string {
const agentName = normalizeAgentName(raw);
if (!agentName) throw new GrokBridgeError("bad_request", 400, "agentName must be a lowercase slug");
return agentName;
}

export async function requirePairing(
get: (id: string) => Promise<GrokPairing | null>,
id: string,
): Promise<GrokPairing> {
const pairing = await get(id);
if (!pairing) throw new GrokBridgeError("not_found", 404, "pairing not found");
return pairing;
}
42 changes: 42 additions & 0 deletions src/grok-bridge/crypto.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,42 @@
import { createHash, randomBytes, timingSafeEqual } from "node:crypto";

export const GROK_BRIDGE_PROTOCOL = "qm-grok-bridge/v1";
export const GROK_BRIDGE_SKILL_REVISION = "qm-grok-bridge/v1";
export const GROK_BRIDGE_QUEUE_CAP = 5;
export const GROK_BRIDGE_DEFAULT_TTL_MS = 30 * 60_000;
export const GROK_BRIDGE_MAX_TTL_MS = 2 * 60 * 60_000;

export function inboundRefFor(pairingId: string): string {
return `grok-bridge:${pairingId}`;
}

export function mintCallbackToken(): string {
return randomBytes(32).toString("base64url");
}

export function hashToken(token: string): string {
return createHash("sha256").update(token).digest("hex");
}

export function tokenMatches(token: string, expectedHash: string): boolean {
const got = hashToken(token);
const a = Buffer.from(got);
const b = Buffer.from(expectedHash);
return a.length === b.length && timingSafeEqual(a, b);
}

export function normalizeAgentName(raw: string): string | undefined {
const name = raw.trim().toLowerCase();
if (!/^[a-z][a-z0-9-]{0,31}$/.test(name)) return undefined;
return name;
}

export function grokDisplayName(agentName: string): string {
return `QM · ${agentName.charAt(0).toUpperCase()}${agentName.slice(1)}`;
}

export function parseBearer(header: string | undefined): string | undefined {
if (!header) return undefined;
const match = /^Bearer\s+(\S+)$/i.exec(header.trim());
return match?.[1];
}
17 changes: 17 additions & 0 deletions src/grok-bridge/durable-vault.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,17 @@
import type { DurableMap } from "../persistence/durable-map.ts";
import type { SecretVault } from "./types.ts";

export interface StoredInboundSecret {
url: string;
bearer: string;
}

export function createDurableSecretVault(backing: DurableMap<StoredInboundSecret>): SecretVault {
return {
async put(id, value) {
await backing.put(id, { url: value.url, bearer: value.bearer });
},
get: (id) => backing.get(id),
delete: (id) => backing.delete(id),
};
}
24 changes: 24 additions & 0 deletions src/grok-bridge/http-outbound.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,24 @@
import type { OutboundPort } from "./types.ts";

const GROK_WEBHOOK_TIMEOUT_MS = 10_000;

export function createHttpOutbound(fetchImpl: typeof fetch = fetch): OutboundPort {
return {
async postJob(url, bearer, envelope) {
try {
const response = await fetchImpl(url, {
method: "POST",
headers: {
authorization: `Bearer ${bearer}`,
"content-type": "application/json",
},
body: JSON.stringify(envelope),
signal: AbortSignal.timeout(GROK_WEBHOOK_TIMEOUT_MS),
});
return { ok: response.ok, status: response.status };
} catch {
return { ok: false, status: 0 };
}
},
};
}
21 changes: 21 additions & 0 deletions src/grok-bridge/job-store.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,21 @@
import { createMemoryMap, type DurableMap } from "../persistence/durable-map.ts";
import type { GrokJob, JobStore } from "./types.ts";

export type { JobStore } from "./types.ts";

export function createJobStore(backing: DurableMap<GrokJob> = createMemoryMap<GrokJob>()): JobStore {
return {
async create(job) {
await backing.put(job.id, job);
return job;
},
get: (id) => backing.get(id),
async save(job) {
await backing.put(job.id, job);
return job;
},
async listByPairing(pairingId) {
return (await backing.all()).filter((job) => job.pairingId === pairingId);
},
};
}
Loading