Skip to content

7. Add a command flow

A “command flow” is the unit of behaviour change in Mnemose: a Zod-typed command that flows through the CQRS spine, produces one or more domain events, and updates read models. This guide adds a hypothetical QuarantineInstance command end-to-end.

Prerequisite: 6. Local development running.

1. Domain schema

Add the command and event Zod schemas:

packages/domain/src/commands/index.ts
export const QuarantineInstanceCommandSchema = TenantScopedCommandBase.extend({
type: z.literal("QuarantineInstance"),
cloudProjectId: z.string().uuid(),
instanceId: z.string(),
reason: z.string().min(1),
});
export type QuarantineInstanceCommand =
z.infer<typeof QuarantineInstanceCommandSchema>;
// Add to AnyCommandSchema discriminated union:
export const AnyCommandSchema = z.discriminatedUnion("type", [
// ...existing
QuarantineInstanceCommandSchema,
]);
packages/domain/src/events/index.ts
export const InstanceQuarantinedSchema = TenantScopedEventBase.extend({
type: z.literal("InstanceQuarantined"),
payload: z.object({
cloudProjectId: z.string().uuid(),
instanceId: z.string(),
reason: z.string(),
}),
});
export type InstanceQuarantinedEvent =
z.infer<typeof InstanceQuarantinedSchema>;
// Add to AnyDomainEventSchema.

Both schemas must be added to their respective discriminated unions or they will not deserialise correctly. Re-export the committed contract catalogs afterwards:

Terminal window
pnpm --filter @mnemose/domain contracts:export

2. Read-model table (only for long-running operations)

If the command has a multi-step lifecycle (pending → running → completed), add a Drizzle table (SQLite/D1 dialect):

packages/db/src/schema/quarantine-jobs.ts
export const quarantineJobs = sqliteTable(
"quarantine_jobs",
{
id: text("id").primaryKey(), // uuid minted by the handler
tenantId: text("tenant_id").notNull(),
cloudProjectId: text("cloud_project_id").notNull(),
instanceId: text("instance_id").notNull(),
status: text("status").$type<"pending" | "running" | "completed" | "failed">().notNull(),
reason: text("reason").notNull(),
createdAt: text("created_at").notNull(),
completedAt: text("completed_at"),
},
);

Export from packages/db/src/schema/index.ts and add a SQL migration under packages/db/src/migrations/ (the D1 migrations are handwritten SQL applied via wrangler d1 migrations).

3. Command handler

Create the handler package:

services/commands/quarantine-instance/src/index.ts
import {
QuarantineInstanceCommandSchema,
InstanceQuarantinedSchema,
} from "@mnemose/domain";
export async function handleQuarantineInstance(
rawPayload: unknown,
deps: { platform: CloudPlatform; db: DbClient },
) {
const cmd = QuarantineInstanceCommandSchema.parse(rawPayload);
// 1. Record the job
const [job] = await deps.db
.insert(quarantineJobs)
.values({
id: crypto.randomUUID(),
tenantId: cmd.tenantId,
cloudProjectId: cmd.cloudProjectId,
instanceId: cmd.instanceId,
status: "running",
reason: cmd.reason,
})
.returning();
// 2. Mint a customer-cloud credential
const project = await loadCloudProject(deps.db, cmd.cloudProjectId);
const credential = await deps.platform.cloudCredentials.mintToken(project.credentialConfig);
try {
// 3. Call the IaaS port
await deps.platform.iaas.compute.applyNetworkTag(
{ credential, projectId: project.projectId, region: project.region },
{ instanceId: cmd.instanceId, tag: "quarantine" },
);
// 4. Write the domain event
await deps.db.insert(domainEvents).values({
aggregateId: cmd.instanceId,
aggregateType: "instance",
eventType: "InstanceQuarantined",
payload: { cloudProjectId: cmd.cloudProjectId, instanceId: cmd.instanceId, reason: cmd.reason },
tenantId: cmd.tenantId,
correlationId: cmd.correlationId,
causationId: cmd.id,
actor: cmd.actor,
});
// 5. Publish for downstream consumers (Cloudflare Queues via EventBus)
await deps.platform.eventBus.publish("mnemose.events", {
type: "InstanceQuarantined",
payload: { cloudProjectId: cmd.cloudProjectId, instanceId: cmd.instanceId, reason: cmd.reason },
// ...envelope
});
await deps.db.update(quarantineJobs)
.set({ status: "completed", completedAt: new Date().toISOString() })
.where(eq(quarantineJobs.id, job.id));
} catch (e) {
await deps.db.update(quarantineJobs).set({ status: "failed" }).where(eq(quarantineJobs.id, job.id));
throw e;
}
}

The handler shape (validate → record job → mint → IaaS → event → publish) is uniform across all command handlers. See services/commands/sync-inventory/src/index.ts for a complete reference.

4. GraphQL mutation

services/gateway/src/schema/quarantine.ts
import { builder } from "./builder.js";
import { QuarantineInstanceCommandSchema } from "@mnemose/domain";
builder.mutationField("quarantineInstance", (t) =>
t.field({
type: "Boolean",
args: {
cloudProjectId: t.arg.string({ required: true }),
instanceId: t.arg.string({ required: true }),
reason: t.arg.string({ required: true }),
},
resolve: async (_root, args, ctx) => {
const cmd = QuarantineInstanceCommandSchema.parse({
type: "QuarantineInstance",
tenantId: ctx.tenantId,
actor: { sub: ctx.actorSub, kind: "user" },
id: crypto.randomUUID(),
correlationId: ctx.correlationId,
timestamp: new Date().toISOString(),
cloudProjectId: args.cloudProjectId,
instanceId: args.instanceId,
reason: args.reason,
});
await ctx.platform.eventBus.publish("mnemose.commands", cmd);
return true;
},
}),
);

Then register the schema:

services/gateway/src/schema/index.ts
import "./quarantine.js";

5. Wire the handler into the queue dispatcher

Handlers are dispatched Workers-edge by the gateway’s queue consumer. Add the new command type to the dispatcher (services/gateway/src/queue-handler.ts and/or the agent’s command subscriber, following how a peer handler like sync-inventory is registered). There is no separate container deploy step — the handler ships inside the Worker bundle.

6. Tests

tests/quarantine-flow.test.ts
import { QuarantineInstanceCommandSchema } from "@mnemose/domain";
test("QuarantineInstance Zod roundtrip", () => {
const cmd = { /* ... */ };
const parsed = QuarantineInstanceCommandSchema.parse(cmd);
expect(parsed.type).toBe("QuarantineInstance");
});

Handler-level tests colocate with the handler package (see peer tests/*.test.ts flow tests).

7. SDL artifact

Regenerate the committed schema and let CI diff it:

Terminal window
pnpm --filter @mnemose/gateway schema:export
git add services/gateway/schema.graphql

Verify end-to-end

Terminal window
pnpm dev # if not already running
# Open GraphiQL at http://localhost:4000/graphql and execute:
mutation { quarantineInstance(cloudProjectId: "...", instanceId: "...", reason: "...") }
# Watch the gateway log: it should consume the command from mnemose.commands,
# run the handler, and emit InstanceQuarantined back through mnemose.events.

Commit discipline

Per CLAUDE.md, this is one atomic commit:

  • Domain schemas
  • DB schema + migration
  • Handler
  • Mutation + SDL update
  • Tests
  • Any docs that need updating (PROJECT.md, docs/api.md)

Never split these — a green commit is one that compiles, type-checks, passes tests, and includes every artifact for the change it claims.

Next

→ 8. Add a cloud adapter