mirror of
https://github.com/mvschwarz/openrig.git
synced 2026-10-02 00:27:21 +08:00
* feat(slack): structured questions with clickable options on human decisions (#193) - `rig queue create --human-questions-file`: 1-4 questions, 2-4 options each (one may be recommended), validated at create; stored by migration 091. - Slack renders each question as a button row with a complete text fallback; a typed thread reply still answers ("Other"). - A block_actions click records that answer on the item; once every question is answered the answers are final, one reply row lands on the seat and the decision resolves like a typed reply. Failed hand-backs are dead-lettered and retried; each click is confirmed in the thread. - App manifest turns on interactivity (Socket Mode, no request URL); setup doc and messaging-the-human skill updated. * fix(slack): own-key answer lookup; require a human destination for questions(#193) * fix(slack): redact secrets in answer acknowledgements (#193)
This commit is contained in:
@@ -76,6 +76,12 @@ So after installing, compare the granted scopes Slack shows for the app with all
|
||||
The app subscribes to messages in public channels it is a member of (`message.channels`) and to
|
||||
mentions of the app (`app_mention`). It does not request direct-message or private-channel access.
|
||||
|
||||
The manifest also turns on **Interactivity**, so the human can answer a decision's structured
|
||||
questions by clicking a button (`rig queue create --human-questions-file`). In Socket Mode the
|
||||
clicks arrive over the same socket, so no request URL is needed. An app created from an older
|
||||
manifest has Interactivity off: turn it on under **Interactivity & Shortcuts**, or the buttons
|
||||
will do nothing. A typed reply in the thread still answers the decision either way.
|
||||
|
||||
## What the connector does with the tokens
|
||||
|
||||
The tokens stay in the env file you created. The connector reads them from that file and uses them
|
||||
|
||||
@@ -427,6 +427,7 @@ export function queueCommand(depsOverride?: QueueDeps): Command {
|
||||
.option("--human-intent <intent>", "decision (default) or update: a quiet informational delivery, never an approval request")
|
||||
.option("--human-detail-file <path>", "One explicitly authored supplemental thread reply; keep the complete action/options in --body-file")
|
||||
.option("--reply-to <qitemId>", "Post this update into an earlier qitem's Slack thread (requires --human-intent update; posts as a new top-level message instead if that thread can't be used, e.g. it is missing or still has an open human decision; --verify reports why)")
|
||||
.option("--human-questions-file <path>", "#193: JSON array of 1-4 questions for a decision, each {id, question, options: [{id, label, recommended?}]} with 2-4 options; Slack shows them as buttons")
|
||||
.option("--evidence-ref <path>", "OPR.0.4.4.19 FR-5: pointer to the durable artifact a human judges (e.g. a PROOF.md path). Required by the daemon when the item is human-routed; optional otherwise.")
|
||||
.option("--host <id>", QUEUE_HOST_OPTION_HELP)
|
||||
.option("--no-nudge", "Suppress the default destination nudge (cold-queue)")
|
||||
@@ -450,6 +451,7 @@ export function queueCommand(depsOverride?: QueueDeps): Command {
|
||||
humanIntent?: string;
|
||||
humanDetailFile?: string;
|
||||
replyTo?: string;
|
||||
humanQuestionsFile?: string;
|
||||
summary?: string;
|
||||
evidenceRef?: string;
|
||||
host?: string;
|
||||
@@ -493,6 +495,26 @@ export function queueCommand(depsOverride?: QueueDeps): Command {
|
||||
process.exitCode = 1;
|
||||
return;
|
||||
}
|
||||
// #193 — read and parse the questions locally too; the daemon validates their shape.
|
||||
let humanQuestions: unknown;
|
||||
if (opts.humanQuestionsFile) {
|
||||
let text: string;
|
||||
try {
|
||||
text = await resolveQueueBody({ bodyFile: opts.humanQuestionsFile });
|
||||
} catch (err) {
|
||||
emitBodyResolveError(err as Error & { fact?: string; consequence?: string; action?: string }, opts.json ?? false);
|
||||
return;
|
||||
}
|
||||
try {
|
||||
humanQuestions = JSON.parse(text);
|
||||
} catch (err) {
|
||||
emitBodyResolveError(Object.assign(new Error(`--human-questions-file ${opts.humanQuestionsFile} is not valid JSON: ${(err as Error).message}`), {
|
||||
consequence: "The queue command did not run; the daemon was not contacted.",
|
||||
action: "Pass a JSON array of questions, each {id, question, options: [{id, label, recommended?}]}.",
|
||||
}), opts.json ?? false);
|
||||
return;
|
||||
}
|
||||
}
|
||||
// OPR.0.4.1.18 (FR-7, warn-then-require grace): a summary SHOULD accompany
|
||||
// every new qitem (it feeds the Story node + helps humans skim). Warn — to
|
||||
// stderr so --json stdout stays clean — but do NOT hard-break existing
|
||||
@@ -547,6 +569,7 @@ export function queueCommand(depsOverride?: QueueDeps): Command {
|
||||
humanIntent: opts.humanIntent,
|
||||
humanDetail: opts.humanDetailFile ? await resolveQueueBody({ bodyFile: opts.humanDetailFile }) : undefined,
|
||||
replyTo: opts.replyTo,
|
||||
humanQuestions,
|
||||
summary: opts.summary,
|
||||
evidenceRef: opts.evidenceRef,
|
||||
priority: opts.priority,
|
||||
|
||||
@@ -141,6 +141,25 @@ describe("rig queue CLI", () => {
|
||||
}
|
||||
});
|
||||
|
||||
it("create --human-questions-file sends the parsed questions; unreadable JSON is refused before any request", async () => {
|
||||
const directory = fs.mkdtempSync(path.join(os.tmpdir(), "queue-human-questions-"));
|
||||
const file = path.join(directory, "questions.json");
|
||||
const questions = [{ id: "db", question: "Which database?", options: [{ id: "pg", label: "Postgres", recommended: true }, { id: "sqlite", label: "SQLite" }] }];
|
||||
fs.writeFileSync(file, JSON.stringify(questions));
|
||||
const broken = path.join(directory, "broken.json");
|
||||
fs.writeFileSync(broken, "{ not json");
|
||||
try {
|
||||
const { deps, calls } = makeDeps();
|
||||
await createProgram({ queueDeps: deps }).parseAsync(["node", "rig", "queue", "create", "--destination", "human-founder@external", "--body", "Pick one.", "--human-intent", "decision", "--human-questions-file", file, "--json"]);
|
||||
expect(calls.find((c) => c.path === "/api/queue/create")?.body).toMatchObject({ humanIntent: "decision", humanQuestions: questions });
|
||||
|
||||
const again = makeDeps();
|
||||
await createProgram({ queueDeps: again.deps }).parseAsync(["node", "rig", "queue", "create", "--destination", "human-founder@external", "--body", "Pick one.", "--human-intent", "decision", "--human-questions-file", broken, "--json"]);
|
||||
expect(process.exitCode).toBe(1);
|
||||
expect(again.calls.find((c) => c.path === "/api/queue/create")).toBeUndefined();
|
||||
} finally { process.exitCode = undefined; fs.rmSync(directory, { recursive: true, force: true }); }
|
||||
});
|
||||
|
||||
// Slice-03 Atom 6b — --body-context snapshot + provenance rule.
|
||||
it("create preserves explicit human intent and authored supplemental file bytes", async () => {
|
||||
const directory = fs.mkdtempSync(path.join(os.tmpdir(), "queue-human-detail-"));
|
||||
|
||||
@@ -92,6 +92,17 @@ the human (a pending human decision, or a row parked on the human), since a
|
||||
reply in that thread would answer the decision. In every such case the
|
||||
`--verify` result says `threaded: false` with the reason.
|
||||
|
||||
A **decision** with a few clear choices can carry `--human-questions-file <path>`:
|
||||
a JSON array of 1–4 questions, each
|
||||
`{"id", "question", "options": [{"id", "label", "recommended"?}]}` with 2–4
|
||||
options (labels up to 75 characters, at most one recommended). Slack shows each
|
||||
question as a row of buttons. Each click records that answer on the item, and
|
||||
the decision resolves once every question has one. You then receive one reply
|
||||
row listing the answers, and the item's `humanAnswers` holds the option ids. The
|
||||
human may instead type a reply in the thread; that resolves the decision as
|
||||
usual, so read the reply rather than assuming an option was picked. Keep the
|
||||
brief complete: the questions add buttons, they do not replace the explanation.
|
||||
|
||||
If an existing agent-owned row must wait for a **decision**, block it on the **new live qitem ID**
|
||||
(`rig queue block <work-id> --on <human-qitem-id> ...`), not on the human address.
|
||||
Completion of the human qitem resumes its dependants. Blocking on the human as
|
||||
|
||||
@@ -92,6 +92,17 @@ the human (a pending human decision, or a row parked on the human), since a
|
||||
reply in that thread would answer the decision. In every such case the
|
||||
`--verify` result says `threaded: false` with the reason.
|
||||
|
||||
A **decision** with a few clear choices can carry `--human-questions-file <path>`:
|
||||
a JSON array of 1–4 questions, each
|
||||
`{"id", "question", "options": [{"id", "label", "recommended"?}]}` with 2–4
|
||||
options (labels up to 75 characters, at most one recommended). Slack shows each
|
||||
question as a row of buttons. Each click records that answer on the item, and
|
||||
the decision resolves once every question has one. You then receive one reply
|
||||
row listing the answers, and the item's `humanAnswers` holds the option ids. The
|
||||
human may instead type a reply in the thread; that resolves the decision as
|
||||
usual, so read the reply rather than assuming an option was picked. Keep the
|
||||
brief complete: the questions add buttons, they do not replace the explanation.
|
||||
|
||||
If an existing agent-owned row must wait for a **decision**, block it on the **new live qitem ID**
|
||||
(`rig queue block <work-id> --on <human-qitem-id> ...`), not on the human address.
|
||||
Completion of the human qitem resumes its dependants. Blocking on the human as
|
||||
|
||||
@@ -93,9 +93,10 @@ import { seatDeliveryGuardSchema } from "./migrations/087_seat_delivery_guard.js
|
||||
import { nodePermissionSelectionsSchema } from "./migrations/088_node_permission_selections.js";
|
||||
import { classificationIdentityProvenanceSchema } from "./migrations/089_classification_identity_provenance.js";
|
||||
import { humanReplyToSchema } from "./migrations/090_human_reply_to.js";
|
||||
import { humanQuestionsSchema } from "./migrations/091_human_questions.js";
|
||||
import type { Migration } from "./migrate.js";
|
||||
|
||||
/** Ordered 001→090 (S02 086/089, S09 087, S03 088, #96 090). */
|
||||
/** Ordered 001→091 (S02 086/089, S09 087, S03 088, #96 090, #193 091). */
|
||||
export const ALL_MIGRATIONS: Migration[] = [
|
||||
coreSchema,
|
||||
bindingsSessionsSchema,
|
||||
@@ -187,4 +188,5 @@ export const ALL_MIGRATIONS: Migration[] = [
|
||||
nodePermissionSelectionsSchema,
|
||||
classificationIdentityProvenanceSchema,
|
||||
humanReplyToSchema,
|
||||
humanQuestionsSchema,
|
||||
];
|
||||
|
||||
@@ -0,0 +1,10 @@
|
||||
import type { Migration } from "../migrate.js";
|
||||
|
||||
// #193 — structured questions on a human decision, and the answers recorded from clicks.
|
||||
export const humanQuestionsSchema: Migration = {
|
||||
name: "091_human_questions.sql",
|
||||
sql: `
|
||||
ALTER TABLE queue_items ADD COLUMN human_questions TEXT;
|
||||
ALTER TABLE queue_items ADD COLUMN human_answers TEXT;
|
||||
`,
|
||||
};
|
||||
@@ -16,6 +16,8 @@ import type { SeenStore, DeadLetterStore, DeadLetterEntry } from "./state-store.
|
||||
import type { InboundQueuePort } from "./queue-access.js";
|
||||
import { createHash } from "node:crypto";
|
||||
import { ADMITTED_EVENT_TYPES } from "./capabilities.js";
|
||||
import { escapeSlackText, parseQuestionAction, redactSecrets } from "./message.js";
|
||||
import { formatHumanAnswers, unansweredQuestions, type RecordHumanAnswerResult } from "../../human-questions.js";
|
||||
|
||||
export interface SlackEvent {
|
||||
type?: string;
|
||||
@@ -32,6 +34,26 @@ export interface SlackEvent {
|
||||
files?: unknown[];
|
||||
}
|
||||
|
||||
/** #193 — the parts of a Socket Mode `block_actions` payload (a button click) we read. */
|
||||
export interface SlackBlockActions {
|
||||
type?: string;
|
||||
user?: { id?: string };
|
||||
channel?: { id?: string };
|
||||
container?: { message_ts?: string; thread_ts?: string; channel_id?: string };
|
||||
message?: { ts?: string; thread_ts?: string };
|
||||
actions?: Array<{ block_id?: string; action_id?: string; value?: string; action_ts?: string }>;
|
||||
}
|
||||
|
||||
export type RecordHumanAnswer = (input: { qitemId: string; actorSession: string; questionId: string; optionId: string }) => RecordHumanAnswerResult;
|
||||
|
||||
/** The root message a click was made on: buttons live on the decision's root post, whose ts is
|
||||
* the thread map's key (a click inside a thread would carry that thread's root instead). */
|
||||
export function clickedRootTs(payload: SlackBlockActions): string | undefined {
|
||||
return payload.container?.thread_ts ?? payload.message?.thread_ts ?? payload.container?.message_ts ?? payload.message?.ts;
|
||||
}
|
||||
|
||||
export type InboundDisposition = "accepted" | "ignored" | "refused" | "dead-lettered" | "handler-failed";
|
||||
|
||||
/**
|
||||
* Loop-safety + non-ingestible ignore. Ingest genuine human messages — text,
|
||||
* and (OPR.0.5.6.2, replacing the T1076 v1 ignore) file-bearing drops: a
|
||||
@@ -97,6 +119,13 @@ export interface InboundDeps {
|
||||
resolveRoute?: (ev: SlackEvent) => { destination: string; tags?: string[]; correlationQitemId?: string };
|
||||
/** Continue an exact human gate through the existing Mission Control resolve primitive. */
|
||||
resolveHumanReply?: (input: { qitemId: string; actorSession: string; decision: string }) => Promise<"resolved" | "already-resolved" | "not-applicable">;
|
||||
/** #193 — record a clicked answer on the decision the clicked message belongs to. */
|
||||
recordHumanAnswer?: RecordHumanAnswer;
|
||||
/** #193 — clicks whose continuation (reply row + resolve) failed; retried with the event
|
||||
* dead-letters. Absent → a failed click is only logged. */
|
||||
actionDeadLetter?: DeadLetterStore<SlackBlockActions>;
|
||||
/** #193 — tell the human in the decision's thread what a click did. Best-effort. */
|
||||
acknowledgeAnswer?: (input: { channel?: string; threadTs: string; text: string }) => Promise<void>;
|
||||
/** OPR.0.5.6.2 — inbound file transfer. Absent with a file-bearing event →
|
||||
* every file is a NAMED failure on the row ("transfer unavailable"), never
|
||||
* a silent drop of message or file. */
|
||||
@@ -255,6 +284,81 @@ export class InboundRouter {
|
||||
return { landed: r.landed, qitemId: r.qitemId, disposition, reason: r.reason, correlationQitemId: r.correlationQitemId, replyResolution: r.replyResolution };
|
||||
}
|
||||
|
||||
/**
|
||||
* #193 — a button click on a decision's structured questions. The clicked message is the
|
||||
* decision's root, so the same thread map a typed reply uses names the decision and its seat
|
||||
* (never the button's own ids, which only pick the question and option). Each click records
|
||||
* one answer; once the set is complete the answers are final, and the continuation lands one
|
||||
* reply row on the seat and resolves the decision, exactly like a typed reply. The reply row's
|
||||
* id is derived from the decision, so a redelivered click or a retry finds the same row. A
|
||||
* failed continuation is dead-lettered and retried like a typed reply that failed to land.
|
||||
*/
|
||||
async routeAction(payload: SlackBlockActions): Promise<{ status: InboundDisposition; reason?: string }> {
|
||||
const r = await this.attemptAction(payload, true);
|
||||
if (r.status === "handler-failed") this.deps.actionDeadLetter?.append(payload, 1);
|
||||
return r;
|
||||
}
|
||||
|
||||
private async attemptAction(payload: SlackBlockActions, live: boolean): Promise<{ status: InboundDisposition; reason?: string }> {
|
||||
const action = payload.actions?.[0];
|
||||
const picked = parseQuestionAction(action?.block_id, action?.action_id);
|
||||
if (!picked) return { status: "ignored", reason: "not-a-question-button" };
|
||||
const who = this.deps.resolveSender(payload.user?.id ?? "");
|
||||
if (!who.admitted) {
|
||||
this.deps.log?.(`click REFUSED — unregistered sender ${payload.user?.id}: ${who.teaching}`);
|
||||
return { status: "refused", reason: "unregistered" };
|
||||
}
|
||||
const rootTs = clickedRootTs(payload);
|
||||
const route = rootTs ? this.deps.resolveRoute?.({ type: "message", thread_ts: rootTs, channel: payload.channel?.id }) : undefined;
|
||||
const qitemId = route?.correlationQitemId;
|
||||
if (!rootTs || !route || !qitemId) return { status: "ignored", reason: "unmapped-message" };
|
||||
const recorded = this.deps.recordHumanAnswer?.({ qitemId, actorSession: who.source, ...picked });
|
||||
if (!recorded || recorded.status !== "recorded") {
|
||||
return { status: "ignored", reason: recorded?.reason ?? "answers-unavailable" };
|
||||
}
|
||||
const acknowledge = async (text: string) => {
|
||||
try {
|
||||
await this.deps.acknowledgeAnswer?.({ channel: payload.channel?.id, threadTs: rootTs, text });
|
||||
} catch (e) {
|
||||
this.deps.log?.(`answer acknowledgement failed qitem=${qitemId}: ${(e as Error).message}`);
|
||||
}
|
||||
};
|
||||
const lines = formatHumanAnswers(recorded.questions, recorded.answers);
|
||||
if (!recorded.complete) {
|
||||
this.deps.log?.(`answer recorded qitem=${qitemId} question=${picked.questionId}`);
|
||||
const [just] = formatHumanAnswers(recorded.questions.filter((q) => q.id === picked.questionId), recorded.answers);
|
||||
const waiting = unansweredQuestions(recorded.questions, recorded.answers).map((q) => q.question);
|
||||
if (live) await acknowledge(escapeSlackText(redactSecrets(`Recorded: ${just}. Still to answer: ${waiting.join("; ")}`)));
|
||||
return { status: "accepted", reason: "answer-recorded" };
|
||||
}
|
||||
let resolution: "resolved" | "already-resolved" | "not-applicable" | undefined;
|
||||
try {
|
||||
await this.deps.queue.createQitem({
|
||||
qitemId: `qitem-slack-answers-${createHash("sha256").update(qitemId).digest("hex").slice(0, 20)}`,
|
||||
source: who.source,
|
||||
destination: route.destination,
|
||||
priority: "routine",
|
||||
tags: [...route.tags ?? ["founder-slack", "inbound"], "human-answer"],
|
||||
summary: `Founder via Slack: answered ${lines.length === 1 ? "1 question" : `${lines.length} questions`}`,
|
||||
body: `${lines.join("\n")}\n\n---\nAnswers (question id → option id): ${JSON.stringify(recorded.answers)}\nIn reply to: ${qitemId} (its humanAnswers field holds the same)\nRouted by openrig slack-inbound (button click).`,
|
||||
});
|
||||
resolution = await this.deps.resolveHumanReply?.({ qitemId, actorSession: who.source, decision: lines.join("; ") });
|
||||
} catch (e) {
|
||||
this.deps.log?.(`answer continuation failed qitem=${qitemId}: ${(e as Error).message}`);
|
||||
if (live) await acknowledge("Your answers are recorded, but handing them back failed. OpenRig will retry.");
|
||||
return { status: "handler-failed", reason: "answer-continuation-failed" };
|
||||
}
|
||||
if (resolution === "not-applicable") {
|
||||
// The reply row reached the seat, but the decision did not close: say so on the log and
|
||||
// the receipt instead of reporting success. Retrying cannot change this outcome.
|
||||
this.deps.log?.(`answers complete but resolve not applicable qitem=${qitemId}`);
|
||||
return { status: "refused", reason: "resolve-not-applicable" };
|
||||
}
|
||||
this.deps.log?.(`answers complete qitem=${qitemId} -> ${route.destination}`);
|
||||
if (resolution === "resolved") await acknowledge(escapeSlackText(redactSecrets(`All answered, sent back: ${lines.join("; ")}`)));
|
||||
return { status: "accepted", reason: "answers-complete" };
|
||||
}
|
||||
|
||||
/**
|
||||
* INTERRUPTION-SAFE retry (item 8): read the durable set NON-destructively,
|
||||
* attempt each, then ATOMICALLY replace the file with only the still-failing
|
||||
@@ -264,7 +368,7 @@ export class InboundRouter {
|
||||
*/
|
||||
async retryDeadLetters(): Promise<{ retried: number; landed: number }> {
|
||||
const entries = this.deps.deadLetter.readAll();
|
||||
if (entries.length === 0) return { retried: 0, landed: 0 };
|
||||
if (entries.length === 0) return this.retryActionDeadLetters();
|
||||
this.deps.log?.(`retrying ${entries.length} dead-letter(s)`);
|
||||
const stillFailing: DeadLetterEntry<SlackEvent>[] = [];
|
||||
let landed = 0;
|
||||
@@ -277,6 +381,24 @@ export class InboundRouter {
|
||||
// reason === "dup" (in-flight) → drop; a concurrent path owns it
|
||||
}
|
||||
this.deps.deadLetter.replaceAll(stillFailing); // atomic; original intact until here
|
||||
const actions = await this.retryActionDeadLetters();
|
||||
return { retried: entries.length + actions.retried, landed: landed + actions.landed };
|
||||
}
|
||||
|
||||
/** #193 — retry dead-lettered clicks. Same interruption-safe shape as the event retry:
|
||||
* read, attempt each, then atomically keep only the ones still failing. */
|
||||
private async retryActionDeadLetters(): Promise<{ retried: number; landed: number }> {
|
||||
const store = this.deps.actionDeadLetter;
|
||||
const entries = store?.readAll() ?? [];
|
||||
if (!store || entries.length === 0) return { retried: 0, landed: 0 };
|
||||
const stillFailing: DeadLetterEntry<SlackBlockActions>[] = [];
|
||||
let landed = 0;
|
||||
for (const e of entries) {
|
||||
const r = await this.attemptAction(e.ev, false);
|
||||
if (r.status === "handler-failed") stillFailing.push({ ev: e.ev, at: e.at, attempts: e.attempts + 1 });
|
||||
else if (r.status === "accepted") landed++;
|
||||
}
|
||||
store.replaceAll(stillFailing);
|
||||
return { retried: entries.length, landed };
|
||||
}
|
||||
}
|
||||
@@ -285,7 +407,7 @@ export interface SocketEnvelope {
|
||||
envelope_id?: string;
|
||||
type?: string;
|
||||
reason?: string;
|
||||
payload?: { event?: SlackEvent };
|
||||
payload?: { event?: SlackEvent } & SlackBlockActions;
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -300,10 +422,15 @@ export async function handleEnvelope(
|
||||
router: InboundRouter,
|
||||
log?: (m: string) => void,
|
||||
onReceived?: () => void,
|
||||
): Promise<{ status: "accepted" | "ignored" | "refused" | "dead-lettered"; reason?: string }> {
|
||||
): Promise<{ status: InboundDisposition; reason?: string }> {
|
||||
if (env.envelope_id) ack(); // fast-ack, unconditional, first
|
||||
onReceived?.(); // diagnostic receipt follows ACK but precedes every handler filter
|
||||
if (env.type === "disconnect") return { status: "ignored", reason: "disconnect" };
|
||||
if (env.type === "interactive") {
|
||||
// #193 — a button click. Other interactive payloads (shortcuts, modals) are not ours.
|
||||
if (env.payload?.type !== "block_actions") return { status: "ignored", reason: "interactive-type" };
|
||||
return router.routeAction(env.payload);
|
||||
}
|
||||
if (env.type !== "events_api") return { status: "ignored", reason: "envelope-type" };
|
||||
const ev = env.payload?.event ?? {};
|
||||
const decision = ingestDecision(ev);
|
||||
|
||||
@@ -75,7 +75,7 @@ export function buildSlackAppManifest(sources: ManifestSources = CANONICAL_MANIF
|
||||
oauth_config: { scopes: { bot: scopes } },
|
||||
settings: {
|
||||
event_subscriptions: { bot_events: events },
|
||||
interactivity: { is_enabled: false },
|
||||
interactivity: { is_enabled: true }, // #193: button clicks arrive over the socket
|
||||
org_deploy_enabled: false,
|
||||
socket_mode_enabled: true,
|
||||
token_rotation_enabled: false,
|
||||
|
||||
@@ -4,6 +4,8 @@
|
||||
// 40,000 truncation; top-level text is the screen-reader/notification fallback)
|
||||
// https://docs.slack.dev/reference/block-kit/blocks/section-block/ (3,000)
|
||||
// https://docs.slack.dev/reference/block-kit/blocks/ (50 blocks)
|
||||
import { MAX_OPTION_LABEL, type HumanQuestion } from "../../human-questions.js";
|
||||
|
||||
export const SLACK_TEXT_CAP = 3900; // Our conservative complete-fallback budget, not Slack’s hard limit.
|
||||
export const SLACK_SECTION_CAP = 3000;
|
||||
/** @deprecated Complete rendering ignores excerpt requests. */
|
||||
@@ -14,6 +16,8 @@ export interface QitemLike {
|
||||
summary?: string | null;
|
||||
body?: string | null;
|
||||
destinationSession?: string | null;
|
||||
/** #193 — structured questions, rendered as one button row per question. */
|
||||
humanQuestions?: readonly HumanQuestion[] | null;
|
||||
}
|
||||
|
||||
/** M1 A5b — an outbound image attachment. A media-bearing OutboundDecision carries these;
|
||||
@@ -92,6 +96,47 @@ export function buildImageBlocks(mediaRefs: readonly SlackMediaRef[] | undefined
|
||||
return blocks;
|
||||
}
|
||||
|
||||
/** #193 — the block_id / action_id prefixes a click carries back. The inbound path parses
|
||||
* exactly these (one producer, one parser: see parseQuestionAction). */
|
||||
export const QUESTION_BLOCK_PREFIX = "or-q:";
|
||||
export const OPTION_ACTION_PREFIX = "or-opt:";
|
||||
const TYPED_REPLY_HINT = "Or reply in this thread with your own answer.";
|
||||
|
||||
/** Parse a clicked button back into its question and option ids; null if it is not ours. */
|
||||
export function parseQuestionAction(blockId: unknown, actionId: unknown): { questionId: string; optionId: string } | null {
|
||||
if (typeof blockId !== "string" || typeof actionId !== "string") return null;
|
||||
if (!blockId.startsWith(QUESTION_BLOCK_PREFIX) || !actionId.startsWith(OPTION_ACTION_PREFIX)) return null;
|
||||
const questionId = blockId.slice(QUESTION_BLOCK_PREFIX.length);
|
||||
const optionId = actionId.slice(OPTION_ACTION_PREFIX.length);
|
||||
return questionId && optionId ? { questionId, optionId } : null;
|
||||
}
|
||||
|
||||
/** #193 — the questions as blocks (a section, then a button row, per question) plus the
|
||||
* complete text they must also appear as in the accessible fallback. */
|
||||
function buildQuestionBlocks(questions: readonly HumanQuestion[]): { blocks: unknown[]; text: string } {
|
||||
const blocks: unknown[] = [];
|
||||
const lines: string[] = [];
|
||||
for (const q of questions) {
|
||||
const question = bounded(`*${inert(q.question)}*`, SLACK_SECTION_CAP, "question");
|
||||
blocks.push({ type: "section", text: { type: "mrkdwn", text: question } });
|
||||
blocks.push({
|
||||
type: "actions",
|
||||
block_id: `${QUESTION_BLOCK_PREFIX}${q.id}`,
|
||||
elements: q.options.map((o) => ({
|
||||
type: "button",
|
||||
action_id: `${OPTION_ACTION_PREFIX}${o.id}`,
|
||||
value: o.id,
|
||||
...(o.recommended ? { style: "primary" } : {}),
|
||||
text: { type: "plain_text", text: bounded(inert(o.label), MAX_OPTION_LABEL, "option label") },
|
||||
})),
|
||||
});
|
||||
lines.push(question, ...q.options.map((o) => `• ${inert(o.label)}${o.recommended ? " (recommended)" : ""}`));
|
||||
}
|
||||
blocks.push({ type: "context", elements: [{ type: "mrkdwn", text: TYPED_REPLY_HINT }] });
|
||||
lines.push(TYPED_REPLY_HINT);
|
||||
return { blocks, text: lines.join("\n") };
|
||||
}
|
||||
|
||||
// Secret-looking patterns we refuse to forward (item 7 defense-in-depth).
|
||||
const SECRET_PATTERNS: RegExp[] = [
|
||||
/xox[baprs]-[A-Za-z0-9-]+/g, // Slack bot/user/app/refresh tokens
|
||||
@@ -199,12 +244,14 @@ export function buildOutboundMessage(q: QitemLike, opts: OutboundMessageOpts): S
|
||||
const imageBlocks = buildImageBlocks(opts.mediaRefs);
|
||||
const attachmentText = imageBlocks.map((b) => `Image: ${(b as { alt_text: string }).alt_text}`).join("\n");
|
||||
const evidence = buildEvidenceLink(opts.evidenceLink);
|
||||
const questionParts = q.humanQuestions?.length ? buildQuestionBlocks(q.humanQuestions) : null;
|
||||
if (opts.extraBlocks?.length) {
|
||||
throw new HumanMessageShapeError("Extra blocks have no complete accessible fallback. Use mediaRefs for images or author supplemental human detail.");
|
||||
}
|
||||
const text = bounded([headline, body, attr, evidence ? evidence.text : null, attachmentText, opts.reconcileMarker].filter(Boolean).join("\n"), SLACK_TEXT_CAP, "complete fallback");
|
||||
const text = bounded([headline, body, questionParts?.text, attr, evidence ? evidence.text : null, attachmentText, opts.reconcileMarker].filter(Boolean).join("\n"), SLACK_TEXT_CAP, "complete fallback");
|
||||
const blocks: unknown[] = [{ type: "section", text: { type: "mrkdwn", text: headline } }];
|
||||
if (body.trim()) blocks.push({ type: "section", text: { type: "mrkdwn", text: body } });
|
||||
if (questionParts) blocks.push(...questionParts.blocks);
|
||||
blocks.push(...imageBlocks);
|
||||
if (evidence) blocks.push(evidence.block);
|
||||
blocks.push({ type: "context", elements: [{ type: "mrkdwn", text: attr }] });
|
||||
|
||||
@@ -32,6 +32,7 @@ export interface OutboundPostPayload {
|
||||
humanDetail?: string | null;
|
||||
/** #96: an update posted into this earlier qitem's thread (subsystem-resolved). */
|
||||
replyTo?: string | null;
|
||||
humanQuestions?: import("../../human-questions.js").HumanQuestion[] | null;
|
||||
summary?: string | null;
|
||||
body?: string | null;
|
||||
destinationSession?: string | null;
|
||||
@@ -136,6 +137,7 @@ function toPayload(q: QueueItem): OutboundPostPayload {
|
||||
humanIntent: q.humanIntent,
|
||||
humanDetail: q.humanDetail,
|
||||
replyTo: q.replyTo,
|
||||
humanQuestions: q.humanQuestions,
|
||||
summary: q.summary,
|
||||
body: q.body,
|
||||
destinationSession: q.destinationSession,
|
||||
|
||||
@@ -7,6 +7,7 @@
|
||||
|
||||
import type { QueueRepository, QueueItem as RepoQueueItem } from "../../queue-repository.js";
|
||||
import { loadHumanRegistry, resolveRegisteredHumanAddress, type LoadResult, type HumanFragment } from "../human-registry.js";
|
||||
import type { HumanQuestion } from "../../human-questions.js";
|
||||
import { ownerNotificationLevelAtLeast, type OwnerNotificationLevel, type QueueTransition } from "../../queue-transition-log.js";
|
||||
|
||||
/** The narrow projection the slack path consumes (shape-compatible with the retired bridge's
|
||||
@@ -21,6 +22,7 @@ export interface QueueItem {
|
||||
humanIntent?: "decision" | "update" | null;
|
||||
humanDetail?: string | null;
|
||||
replyTo?: string | null;
|
||||
humanQuestions?: HumanQuestion[] | null;
|
||||
summary?: string | null;
|
||||
body?: string | null;
|
||||
evidenceRef?: string | null;
|
||||
@@ -89,6 +91,7 @@ function project(q: RepoQueueItem, transition: QueueTransition, entities: readon
|
||||
humanIntent: q.humanIntent,
|
||||
humanDetail: q.humanDetail,
|
||||
replyTo: q.replyTo ?? null,
|
||||
humanQuestions: q.humanQuestions ?? null,
|
||||
summary: (r.summary as string | null) ?? null,
|
||||
body: (r.body as string | null) ?? null,
|
||||
evidenceRef: (r.evidenceRef as string | null) ?? null,
|
||||
|
||||
@@ -175,6 +175,7 @@ function deliverSinglePart(opts: SubsystemSlackDeliveryOpts, markEpisode = true)
|
||||
qitemId: q.qitemId ?? decision.decisionId,
|
||||
summary: q.summary,
|
||||
body: q.body,
|
||||
humanQuestions: q.humanQuestions,
|
||||
destinationSession: q.destinationSession ?? decision.entityBindingRef,
|
||||
},
|
||||
{
|
||||
@@ -386,7 +387,7 @@ export function subsystemSlackDeliver(opts: SubsystemSlackDeliveryOpts): Subsyst
|
||||
const parts = q.humanDetail
|
||||
? [
|
||||
{ ...q, humanDetail: undefined, body: `${q.body ?? ""}\n\nSupplemental detail follows in this thread.` },
|
||||
{ ...q, humanDetail: undefined, summary: `Supplemental detail: ${q.summary ?? ""}`, body: q.humanDetail, media: [], evidenceRef: null },
|
||||
{ ...q, humanDetail: undefined, humanQuestions: undefined, summary: `Supplemental detail: ${q.summary ?? ""}`, body: q.humanDetail, media: [], evidenceRef: null },
|
||||
]
|
||||
: [q];
|
||||
const partId = (index: number) => parts.length === 1 ? decision.decisionId : `${decision.decisionId}:part:${index + 1}`;
|
||||
|
||||
@@ -16,14 +16,14 @@ import { channelStateDigest } from "../channel-operations.js";
|
||||
import path from "node:path";
|
||||
import fs from "node:fs";
|
||||
import { buildInProcessWire, type GatewayWire, type SubsystemDeliverFn } from "../gateway-subsystem.js";
|
||||
import { downloadPrivateFile } from "./slack-api.js";
|
||||
import { downloadPrivateFile, postChatMessage } from "./slack-api.js";
|
||||
import { loadConfig } from "./config.js";
|
||||
import { resolveSecret } from "./secrets.js";
|
||||
import { SeenStore, DeadLetterStore, InboundReceiptStore } from "./state-store.js";
|
||||
import { makeQueuePorts } from "./queue-access.js";
|
||||
import { SlackOutboundDriver, OUTBOUND_OP, type OutboundPostPayload } from "./outbound-driver.js";
|
||||
import { subsystemSlackDeliver } from "./slack-delivery.js";
|
||||
import { InboundRouter, type SlackEvent, type InboundFilePort, type InboundFileResult, type StoredInboundFile, type FailedInboundFile } from "./inbound.js";
|
||||
import { InboundRouter, type SlackEvent, type SlackBlockActions, type InboundFilePort, type InboundFileResult, type StoredInboundFile, type FailedInboundFile } from "./inbound.js";
|
||||
import { makeInboundSenderResolver, type RegistrySurface } from "./inbound-admission.js";
|
||||
import { ThreadSeatMap, formatPostedStamp } from "./thread-seat-map.js";
|
||||
import { makeThreadRouteResolver } from "./thread-routing.js";
|
||||
@@ -571,6 +571,19 @@ export function buildSlackGatewayWire(opts: SlackWireOpts): GatewayWire {
|
||||
// or human-initiated → the configured orchestrator slot as an unrouted-signal row.
|
||||
resolveRoute: makeThreadRouteResolver({ map: threadMap, unroutedDestination: cfg.inboundDestination, log }),
|
||||
resolveHumanReply: opts.resolveHumanReply,
|
||||
// #193 — a button click records its answer on the decision the clicked root belongs to;
|
||||
// a failed hand-back is retried with the event dead-letters, and each click is confirmed
|
||||
// in the decision's thread (a bot post, so inbound never ingests it).
|
||||
recordHumanAnswer: (input) => opts.queueRepo.recordHumanAnswer(input),
|
||||
actionDeadLetter: new DeadLetterStore<SlackBlockActions>(path.join(stateDir(opts.home), "slack-inbound-action-deadletter.jsonl")),
|
||||
...(bot ? {
|
||||
acknowledgeAnswer: async ({ channel, threadTs, text }: { channel?: string; threadTs: string; text: string }) => {
|
||||
const target = channel ?? cfg.channel;
|
||||
if (!target) return;
|
||||
const r = await postChatMessage(bot, { channel: target, thread_ts: threadTs, text }, opts.fetchImpl);
|
||||
if (!r.ok) log(`answer acknowledgement not posted thread=${threadTs}: ${r.error}`);
|
||||
},
|
||||
} : {}),
|
||||
// OPR.0.5.6.2 — inbound file transfer: wired only when the bot token exists
|
||||
// (downloads need `files:read`); absent → the router's own named-failure
|
||||
// arm keeps failure honest. Media lives beside the gateway's other durable
|
||||
|
||||
@@ -0,0 +1,93 @@
|
||||
// #193 — structured questions a decision may carry: 1–4 questions, each with 2–4 options the
|
||||
// human can click in Slack. The validated shape is stored on the item; each click records one
|
||||
// answer, and the decision resolves once every question has one. Limits follow Slack Block Kit:
|
||||
// a button's plain_text is at most 75 characters.
|
||||
|
||||
export interface HumanQuestionOption {
|
||||
id: string;
|
||||
label: string;
|
||||
recommended?: boolean;
|
||||
}
|
||||
|
||||
export interface HumanQuestion {
|
||||
id: string;
|
||||
question: string;
|
||||
options: HumanQuestionOption[];
|
||||
}
|
||||
|
||||
/** questionId → chosen optionId. */
|
||||
export type HumanAnswers = Record<string, string>;
|
||||
|
||||
/** The outcome of recording one clicked answer (QueueRepository.recordHumanAnswer). */
|
||||
export type RecordHumanAnswerResult =
|
||||
| { status: "recorded"; answers: HumanAnswers; complete: boolean; questions: HumanQuestion[] }
|
||||
| { status: "not-applicable"; reason: string };
|
||||
|
||||
export const MAX_HUMAN_QUESTIONS = 4;
|
||||
export const MIN_OPTIONS = 2;
|
||||
export const MAX_OPTIONS = 4;
|
||||
export const MAX_OPTION_LABEL = 75; // Slack button text limit
|
||||
export const MAX_QUESTION_TEXT = 500;
|
||||
const ID_PATTERN = /^[A-Za-z0-9_-]{1,40}$/;
|
||||
const ID_RULE = 'id must be 1–40 letters, digits, "-" or "_".';
|
||||
|
||||
export type HumanQuestionsParse = { ok: true; questions: HumanQuestion[] } | { ok: false; error: string };
|
||||
|
||||
/** Validate an authored questions value. Returns a normalized copy (only known fields kept). */
|
||||
export function parseHumanQuestions(value: unknown): HumanQuestionsParse {
|
||||
const fail = (error: string): HumanQuestionsParse => ({ ok: false, error });
|
||||
if (!Array.isArray(value)) return fail("humanQuestions must be an array of questions.");
|
||||
if (value.length < 1 || value.length > MAX_HUMAN_QUESTIONS) return fail(`Ask 1 to ${MAX_HUMAN_QUESTIONS} questions (got ${value.length}).`);
|
||||
const questionIds = new Set<string>();
|
||||
const questions: HumanQuestion[] = [];
|
||||
for (const [questionIndex, raw] of value.entries()) {
|
||||
const q = (raw ?? {}) as Record<string, unknown>;
|
||||
const where = `question ${questionIndex + 1}`;
|
||||
if (typeof q.id !== "string" || !ID_PATTERN.test(q.id)) return fail(`${where}: ${ID_RULE}`);
|
||||
if (questionIds.has(q.id)) return fail(`${where}: duplicate question id "${q.id}".`);
|
||||
questionIds.add(q.id);
|
||||
if (typeof q.question !== "string" || !q.question.trim()) return fail(`${where}: question text is required.`);
|
||||
if (q.question.length > MAX_QUESTION_TEXT) return fail(`${where}: question text is over ${MAX_QUESTION_TEXT} characters.`);
|
||||
if (!Array.isArray(q.options) || q.options.length < MIN_OPTIONS || q.options.length > MAX_OPTIONS) {
|
||||
return fail(`${where}: give ${MIN_OPTIONS} to ${MAX_OPTIONS} options.`);
|
||||
}
|
||||
const optionIds = new Set<string>();
|
||||
const options: HumanQuestionOption[] = [];
|
||||
let recommended = 0;
|
||||
for (const [optionIndex, rawOption] of q.options.entries()) {
|
||||
const o = (rawOption ?? {}) as Record<string, unknown>;
|
||||
const optionWhere = `${where} option ${optionIndex + 1}`;
|
||||
if (typeof o.id !== "string" || !ID_PATTERN.test(o.id)) return fail(`${optionWhere}: ${ID_RULE}`);
|
||||
if (optionIds.has(o.id)) return fail(`${optionWhere}: duplicate option id "${o.id}".`);
|
||||
optionIds.add(o.id);
|
||||
if (typeof o.label !== "string" || !o.label.trim()) return fail(`${optionWhere}: label is required.`);
|
||||
if (o.label.length > MAX_OPTION_LABEL) return fail(`${optionWhere}: label is over Slack's ${MAX_OPTION_LABEL}-character button limit.`);
|
||||
if (o.recommended != null && typeof o.recommended !== "boolean") return fail(`${optionWhere}: recommended must be true or false.`);
|
||||
if (o.recommended) recommended++;
|
||||
options.push(o.recommended ? { id: o.id, label: o.label, recommended: true } : { id: o.id, label: o.label });
|
||||
}
|
||||
if (recommended > 1) return fail(`${where}: mark at most one option as recommended.`);
|
||||
questions.push({ id: q.id, question: q.question, options });
|
||||
}
|
||||
return { ok: true, questions };
|
||||
}
|
||||
|
||||
/** The recorded option for a question. Own keys only: a question id like "constructor" must
|
||||
* not read the inherited Object.prototype member as an answer. */
|
||||
function ownAnswer(answers: HumanAnswers, questionId: string): string | undefined {
|
||||
return Object.prototype.hasOwnProperty.call(answers, questionId) ? answers[questionId] : undefined;
|
||||
}
|
||||
|
||||
/** One "question: chosen label" line per answered question, in question order. */
|
||||
export function formatHumanAnswers(questions: readonly HumanQuestion[], answers: HumanAnswers): string[] {
|
||||
return questions.flatMap((q) => {
|
||||
const optionId = ownAnswer(answers, q.id);
|
||||
if (optionId == null) return [];
|
||||
return [`${q.question}: ${q.options.find((o) => o.id === optionId)?.label ?? optionId}`];
|
||||
});
|
||||
}
|
||||
|
||||
/** The questions still waiting for an answer. */
|
||||
export function unansweredQuestions(questions: readonly HumanQuestion[], answers: HumanAnswers): HumanQuestion[] {
|
||||
return questions.filter((q) => ownAnswer(answers, q.id) == null);
|
||||
}
|
||||
@@ -27,6 +27,7 @@ import {
|
||||
} from "./queue-wake-repository.js";
|
||||
import { WatchdogJobsRepository } from "./watchdog-jobs-repository.js";
|
||||
import { armQueueWait, backOffQueueWait, refreshQueueWaits, evaluateQueueWait, retargetQueueWait, isQueueWait } from "./queue-wait-backoff.js";
|
||||
import { parseHumanQuestions, unansweredQuestions, type HumanQuestion, type HumanAnswers, type RecordHumanAnswerResult } from "./human-questions.js";
|
||||
|
||||
export const QUEUE_STATES = [
|
||||
"pending",
|
||||
@@ -167,6 +168,10 @@ export interface QueueItem {
|
||||
/** #96 — set on full reads of a replyTo row: why delivery posted top-level instead
|
||||
* (e.g. `root-missing`), or null when it threaded / has not been delivered. */
|
||||
replyToFallback?: string | null;
|
||||
/** #193 — structured questions on a decision (clickable options in Slack). */
|
||||
humanQuestions?: HumanQuestion[] | null;
|
||||
/** #193 — answers recorded from clicks, questionId → optionId; null until the first click. */
|
||||
humanAnswers?: HumanAnswers | null;
|
||||
/** Short human-readable subject; null for callers that omit it. */
|
||||
summary: string | null;
|
||||
/** OPR.0.4.4.19 FR-5 — pointer to the durable artifact a human judges
|
||||
@@ -210,6 +215,8 @@ interface QueueItemRow {
|
||||
human_intent?: "decision" | "update" | null;
|
||||
human_detail?: string | null;
|
||||
reply_to?: string | null;
|
||||
human_questions?: string | null;
|
||||
human_answers?: string | null;
|
||||
summary: string | null;
|
||||
evidence_ref: string | null;
|
||||
closure_reason: string | null;
|
||||
@@ -271,6 +278,8 @@ export interface QueueCreateInput {
|
||||
* that thread can't be used (e.g. a live human gate, see hasLiveHumanGate), it posts
|
||||
* top-level and the row records why. */
|
||||
replyTo?: string | null;
|
||||
/** #193 — 1–4 structured questions; accepted only with humanIntent "decision". */
|
||||
humanQuestions?: HumanQuestion[] | null;
|
||||
summary?: string | null;
|
||||
/** OPR.0.4.4.19 FR-5 — optional durable-artifact pointer. Persisted when
|
||||
* present; required at the domain layer only for human-routed items. */
|
||||
@@ -665,6 +674,7 @@ export class QueueRepository {
|
||||
private readonly hasSummaryColumn: boolean;
|
||||
private readonly hasHumanIntentColumn: boolean;
|
||||
private readonly hasReplyToColumn: boolean;
|
||||
private readonly hasHumanQuestionsColumn: boolean;
|
||||
private readonly hasEvidenceRefColumn: boolean;
|
||||
private readonly hasMintingGenColumn: boolean;
|
||||
private readonly hasClaimedGenColumn: boolean;
|
||||
@@ -725,6 +735,7 @@ export class QueueRepository {
|
||||
this.hasSummaryColumn = detectQueueColumn(db, "summary");
|
||||
this.hasHumanIntentColumn = detectQueueColumn(db, "human_intent");
|
||||
this.hasReplyToColumn = detectQueueColumn(db, "reply_to");
|
||||
this.hasHumanQuestionsColumn = detectQueueColumn(db, "human_questions");
|
||||
this.hasEvidenceRefColumn = detectQueueColumn(db, "evidence_ref");
|
||||
this.hasQueueTransitionsTable = detectTable(db, "queue_transitions");
|
||||
const transitionColumns = this.hasQueueTransitionsTable
|
||||
@@ -1488,6 +1499,20 @@ export class QueueRepository {
|
||||
if (!this.hasHumanIntentColumn) throw new QueueRepositoryError("invalid_human_notification", "Human notification fields require the current queue schema; they were not saved.");
|
||||
}
|
||||
if (input.replyTo != null) this.validateReplyTo(input.replyTo, input.humanIntent);
|
||||
let humanQuestions: HumanQuestion[] | null = null;
|
||||
if (input.humanQuestions != null) {
|
||||
if (!isHumanSeatSessionRef(input.destinationSession)) {
|
||||
throw new QueueRepositoryError("invalid_human_questions", "humanQuestions require a human destination: only a human can click the options.");
|
||||
}
|
||||
// Omitted intent is a decision (legacy behavior), so it may carry questions too.
|
||||
if (input.humanIntent === "update") {
|
||||
throw new QueueRepositoryError("invalid_human_questions", "humanQuestions are refused on an update: they ask the human to decide. Use humanIntent decision.");
|
||||
}
|
||||
const parsed = parseHumanQuestions(input.humanQuestions);
|
||||
if (!parsed.ok) throw new QueueRepositoryError("invalid_human_questions", parsed.error);
|
||||
if (!this.hasHumanQuestionsColumn) throw new QueueRepositoryError("invalid_human_questions", "humanQuestions require the current queue schema; they were not saved.");
|
||||
humanQuestions = parsed.questions;
|
||||
}
|
||||
const id = input.qitemId ?? newQitemId();
|
||||
const ts = new Date().toISOString();
|
||||
const priority = input.priority ?? "routine";
|
||||
@@ -1527,6 +1552,9 @@ export class QueueRepository {
|
||||
if (input.replyTo != null) {
|
||||
this.db.prepare("UPDATE queue_items SET reply_to = ? WHERE qitem_id = ?").run(input.replyTo, id);
|
||||
}
|
||||
if (humanQuestions) {
|
||||
this.db.prepare("UPDATE queue_items SET human_questions = ? WHERE qitem_id = ?").run(JSON.stringify(humanQuestions), id);
|
||||
}
|
||||
this.persistMintingGeneration(id, input.sourceSession);
|
||||
const notification = this.classifyOwnerNotification({
|
||||
action: "create",
|
||||
@@ -3011,6 +3039,34 @@ export class QueueRepository {
|
||||
for (const event of events) this.eventBus.notifySubscribers(event);
|
||||
}
|
||||
|
||||
/**
|
||||
* #193 — record one clicked answer on a pending decision that carries structured questions.
|
||||
* Only the decision's own human may answer, and only while it is pending. A click may change
|
||||
* an earlier answer until the set is complete; from then on the answers are FINAL, and any
|
||||
* further click returns them unchanged so the caller can retry the continuation (reply row +
|
||||
* resolve) with exactly what the seat will read. Anything else is not-applicable (a forged
|
||||
* id, a stranger, a decision already resolved), never an error the socket must retry.
|
||||
*/
|
||||
recordHumanAnswer(input: { qitemId: string; actorSession: string; questionId: string; optionId: string }): RecordHumanAnswerResult {
|
||||
if (!this.hasHumanQuestionsColumn) return { status: "not-applicable", reason: "schema" };
|
||||
return this.db.transaction(() => {
|
||||
const item = this.getById(input.qitemId);
|
||||
if (!item?.humanQuestions?.length) return { status: "not-applicable" as const, reason: "no-questions" };
|
||||
if (item.state !== "pending") return { status: "not-applicable" as const, reason: `state-${item.state}` };
|
||||
if (item.destinationSession !== input.actorSession) return { status: "not-applicable" as const, reason: "not-the-asked-human" };
|
||||
const question = item.humanQuestions.find((q) => q.id === input.questionId);
|
||||
if (!question?.options.some((o) => o.id === input.optionId)) return { status: "not-applicable" as const, reason: "unknown-option" };
|
||||
const recorded = item.humanAnswers ?? {};
|
||||
if (unansweredQuestions(item.humanQuestions, recorded).length === 0) {
|
||||
return { status: "recorded" as const, answers: recorded, complete: true, questions: item.humanQuestions };
|
||||
}
|
||||
const answers: HumanAnswers = { ...(item.humanAnswers ?? {}), [input.questionId]: input.optionId };
|
||||
this.db.prepare("UPDATE queue_items SET human_answers = ? WHERE qitem_id = ?").run(JSON.stringify(answers), input.qitemId);
|
||||
const complete = unansweredQuestions(item.humanQuestions, answers).length === 0;
|
||||
return { status: "recorded" as const, answers, complete, questions: item.humanQuestions };
|
||||
})();
|
||||
}
|
||||
|
||||
getById(qitemId: string): QueueItem | null {
|
||||
const row = this.db
|
||||
.prepare("SELECT * FROM queue_items WHERE qitem_id = ?")
|
||||
@@ -3647,6 +3703,8 @@ export class QueueRepository {
|
||||
humanIntent: row.human_intent ?? null,
|
||||
humanDetail: row.human_detail ?? null,
|
||||
replyTo: row.reply_to ?? null,
|
||||
humanQuestions: row.human_questions ? (JSON.parse(row.human_questions) as HumanQuestion[]) : null,
|
||||
humanAnswers: row.human_answers ? (JSON.parse(row.human_answers) as HumanAnswers) : null,
|
||||
closureReason: row.closure_reason as ClosureReason | null,
|
||||
closureTarget: row.closure_target,
|
||||
closureRequiredAt: row.closure_required_at,
|
||||
|
||||
@@ -21,6 +21,7 @@ import { LOCAL_HOST_ID } from "../domain/hosts/fanout-contract.js";
|
||||
import { remoteJsonRequest } from "../domain/hosts/remote-daemon-http.js";
|
||||
import type { SettingsStore } from "../domain/user-settings/settings-store.js";
|
||||
import { deriveCurrentWork } from "../domain/current-work.js";
|
||||
import type { HumanQuestion } from "../domain/human-questions.js";
|
||||
|
||||
/**
|
||||
* Coordination L3 — Queue HTTP routes (PL-004 Phase A).
|
||||
@@ -145,6 +146,7 @@ export function queueRoutes(): Hono {
|
||||
// #96: --reply-to refusals are named client errors.
|
||||
: err.code === "reply_to_requires_update" ? 400
|
||||
: err.code === "reply_to_not_found" ? 400
|
||||
: err.code === "invalid_human_questions" ? 400
|
||||
// OPR.0.5.1 slice-51-06 D2: summary/evidence_ref on a non-park transition — a client
|
||||
// input error surfaced as a structured 400 (the daemon rejects before any mutation).
|
||||
: err.code === "summary_evidence_not_persistable" ? 400
|
||||
@@ -404,6 +406,7 @@ export function queueRoutes(): Hono {
|
||||
humanIntent?: "decision" | "update" | null;
|
||||
humanDetail?: string | null;
|
||||
replyTo?: string | null;
|
||||
humanQuestions?: HumanQuestion[] | null; // shape validated by the repository (invalid_human_questions)
|
||||
summary?: string | null;
|
||||
evidenceRef?: string | null;
|
||||
nudge?: boolean;
|
||||
@@ -473,6 +476,7 @@ export function queueRoutes(): Hono {
|
||||
humanIntent: body.humanIntent,
|
||||
humanDetail: body.humanDetail,
|
||||
replyTo: body.replyTo,
|
||||
humanQuestions: body.humanQuestions,
|
||||
summary: body.summary,
|
||||
evidenceRef: body.evidenceRef,
|
||||
nudge: (body as { nudge?: boolean }).nudge,
|
||||
|
||||
@@ -55,6 +55,7 @@ import { inventoryEventIndexesSchema } from "../../src/db/migrations/084_invento
|
||||
import { rigClaudeManagedBlockFileSchema } from "../../src/db/migrations/085_rig_claude_managed_block_file.js";
|
||||
import { nodePermissionSelectionsSchema } from "../../src/db/migrations/088_node_permission_selections.js";
|
||||
import { humanReplyToSchema } from "../../src/db/migrations/090_human_reply_to.js";
|
||||
import { humanQuestionsSchema } from "../../src/db/migrations/091_human_questions.js";
|
||||
import { rigPolicySchema } from "../../src/db/migrations/041_rig_policy.js";
|
||||
import { rigArchiveSchema } from "../../src/db/migrations/042_rig_archive.js";
|
||||
import { resumeProvenanceSchema } from "../../src/db/migrations/043_resume_provenance.js";
|
||||
@@ -117,7 +118,7 @@ import fs from "node:fs";
|
||||
|
||||
/** Seam B R6: the canonical full-fixture migration list, exported so file-backed
|
||||
* DB-reopen tests migrate IDENTICALLY to createFullTestDb. */
|
||||
export const migrationsForFullTestDb = [coreSchema, bindingsSessionsSchema, eventsSchema, snapshotsSchema, checkpointsSchema, resumeMetadataSchema, nodeSpecFieldsSchema, packagesSchema, installJournalSchema, journalSeqSchema, bootstrapSchema, discoverySchema, discoveryFkFix, agentspecRebootSchema, startupContextSchema, chatMessagesSchema, podNamespaceSchema, contextUsageSchema, externalCliAttachmentSchema, rigServicesSchema, seatHandoverObservabilitySchema, nodeCodexConfigProfileSchema, nodePermissionPolicySchema, rigPermissionPolicySchema, nodePolicyProvenanceSchema, rigPolicyProvenanceSchema, streamItemsSchema, queueItemsSchema, queueTransitionsSchema, rigPolicySchema, rigArchiveSchema, resumeProvenanceSchema, resumeVerificationSchema, seatIdentityVerdictsSchema, selfHostIdentitySchema, occupantTenuresSchema, daemonLifecycleSchema, watchdogJobsSchema, occupantGenerationStampsSchema, projectionManifestSchema, watchdogTargetGenerationSchema, appliedLaunchObservationsSchema, appliedLaunchObservationInvalidationsSchema, threadSeatMapSchema, queueTransitionWakesSchema, nodeSessionSourceSchema, scopedOperatingPostureSchema, humanNotificationIntentSchema, reviewReadIndexesSchema, inventoryEventIndexesSchema, rigClaudeManagedBlockFileSchema, nodePermissionSelectionsSchema, humanReplyToSchema];
|
||||
export const migrationsForFullTestDb = [coreSchema, bindingsSessionsSchema, eventsSchema, snapshotsSchema, checkpointsSchema, resumeMetadataSchema, nodeSpecFieldsSchema, packagesSchema, installJournalSchema, journalSeqSchema, bootstrapSchema, discoverySchema, discoveryFkFix, agentspecRebootSchema, startupContextSchema, chatMessagesSchema, podNamespaceSchema, contextUsageSchema, externalCliAttachmentSchema, rigServicesSchema, seatHandoverObservabilitySchema, nodeCodexConfigProfileSchema, nodePermissionPolicySchema, rigPermissionPolicySchema, nodePolicyProvenanceSchema, rigPolicyProvenanceSchema, streamItemsSchema, queueItemsSchema, queueTransitionsSchema, rigPolicySchema, rigArchiveSchema, resumeProvenanceSchema, resumeVerificationSchema, seatIdentityVerdictsSchema, selfHostIdentitySchema, occupantTenuresSchema, daemonLifecycleSchema, watchdogJobsSchema, occupantGenerationStampsSchema, projectionManifestSchema, watchdogTargetGenerationSchema, appliedLaunchObservationsSchema, appliedLaunchObservationInvalidationsSchema, threadSeatMapSchema, queueTransitionWakesSchema, nodeSessionSourceSchema, scopedOperatingPostureSchema, humanNotificationIntentSchema, reviewReadIndexesSchema, inventoryEventIndexesSchema, rigClaudeManagedBlockFileSchema, nodePermissionSelectionsSchema, humanReplyToSchema, humanQuestionsSchema];
|
||||
|
||||
/**
|
||||
* P24 — the DECLARED exclusions for {@link migrationsForFullTestDb}. That list is deliberately a
|
||||
|
||||
@@ -48,10 +48,11 @@ describe("shipped Slack app manifest — canonical sources", () => {
|
||||
}
|
||||
});
|
||||
|
||||
it("is a Socket Mode app with no request URL, interactivity or org-wide deploy", () => {
|
||||
it("is a Socket Mode app with interactivity (button clicks, #193) but no request URL or org-wide deploy", () => {
|
||||
const s = bundle.manifest.settings;
|
||||
expect(s.socket_mode_enabled).toBe(true);
|
||||
expect(s.interactivity.is_enabled).toBe(false);
|
||||
// In Socket Mode, block_actions arrive over the socket: interactivity needs no request URL.
|
||||
expect(s.interactivity.is_enabled).toBe(true);
|
||||
expect(s.org_deploy_enabled).toBe(false);
|
||||
expect(JSON.stringify(bundle.manifest)).not.toMatch(/request_url|redirect_url/);
|
||||
});
|
||||
|
||||
@@ -0,0 +1,328 @@
|
||||
// #193 — a decision may carry 1–4 structured questions, each with 2–4 clickable options.
|
||||
// Slack renders them as buttons; a click (block_actions over Socket Mode) records that
|
||||
// question's answer on the same item, and the decision resolves once every question has one.
|
||||
// A typed reply in the thread still resolves the decision as before ("Other").
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
|
||||
import { mkdtempSync, rmSync, writeFileSync } from "node:fs";
|
||||
import { tmpdir } from "node:os";
|
||||
import { join } from "node:path";
|
||||
import { Hono } from "hono";
|
||||
import { queueRoutes } from "../src/routes/queue.js";
|
||||
import { createDb } from "../src/db/connection.js";
|
||||
import { migrate } from "../src/db/migrate.js";
|
||||
import { ALL_MIGRATIONS } from "../src/db/all-migrations.js";
|
||||
import { EventBus } from "../src/domain/event-bus.js";
|
||||
import { QueueRepository } from "../src/domain/queue-repository.js";
|
||||
import { buildOutboundMessage } from "../src/domain/gateway/slack/message.js";
|
||||
import { makeQueuePorts } from "../src/domain/gateway/slack/queue-access.js";
|
||||
import { buildSlackGatewayWire, makeHumanReplyResolver } from "../src/domain/gateway/slack/slack-subsystem.js";
|
||||
import type { WsLike } from "../src/domain/gateway/slack/socket-inbound.js";
|
||||
import { DEFAULT_CONFIG, saveConfig } from "../src/domain/gateway/slack/config.js";
|
||||
import { resolveSlackHandle } from "../src/domain/gateway/human-registry.js";
|
||||
import { MissionControlActionLog } from "../src/domain/mission-control/mission-control-action-log.js";
|
||||
import { MissionControlWriteContract } from "../src/domain/mission-control/mission-control-write-contract.js";
|
||||
import { InboundReceiptStore } from "../src/domain/gateway/slack/state-store.js";
|
||||
|
||||
const human = "human-founder@external";
|
||||
const registry = { ok: true as const, entities: [{ entityId: "human-founder", class: "human" as const, displayName: "Founder", address: human, connectorBindings: [{ kind: "slack" as const, connectorRef: "primary", secretsRef: "env:SLACK_BOT_TOKEN", role: "primary" as const, handle: "UFOUNDER" }], prefs: { deliveryClass: "A" as const } }] };
|
||||
const request = { sourceSession: "author@rig", destinationSession: human, summary: "Two quick questions", body: "Pick one for each; my recommendation is marked.", evidenceRef: "/private/proof.md", nudge: false };
|
||||
const questions = [
|
||||
{ id: "db", question: "Which database?", options: [{ id: "pg", label: "Postgres", recommended: true }, { id: "sqlite", label: "SQLite" }] },
|
||||
{ id: "ship", question: "Ship this week?", options: [{ id: "yes", label: "Yes" }, { id: "no", label: "No" }] },
|
||||
];
|
||||
|
||||
describe("structured human questions (#193)", () => {
|
||||
let home: string;
|
||||
let db: ReturnType<typeof createDb>;
|
||||
let bus: EventBus;
|
||||
let repo: QueueRepository;
|
||||
const stops: Array<() => void> = [];
|
||||
beforeEach(() => {
|
||||
home = mkdtempSync(join(tmpdir(), "human-questions-"));
|
||||
db = createDb(); migrate(db, ALL_MIGRATIONS);
|
||||
bus = new EventBus(db);
|
||||
repo = new QueueRepository(db, bus, { loadHumanRegistry: () => registry });
|
||||
});
|
||||
afterEach(() => { for (const stop of stops.splice(0)) stop(); db.close(); rmSync(home, { recursive: true, force: true }); });
|
||||
|
||||
describe("create", () => {
|
||||
it("records the questions on a decision item", async () => {
|
||||
const created = await repo.create({ ...request, humanIntent: "decision", humanQuestions: questions });
|
||||
expect(repo.getById(created.qitemId)).toMatchObject({ humanQuestions: questions, humanAnswers: null });
|
||||
});
|
||||
|
||||
it("accepts questions over HTTP, and answers a bad shape with a 400 that names the problem", async () => {
|
||||
const app = new Hono();
|
||||
app.use("*", async (c, next) => { (c.set as (k: string, v: unknown) => void)("queueRepo", repo); await next(); });
|
||||
app.route("/api/queue", queueRoutes());
|
||||
const create = (value: unknown) => app.request("/api/queue/create", { method: "POST", headers: { "Content-Type": "application/json", "X-OpenRig-Session": request.sourceSession }, body: JSON.stringify(value) });
|
||||
const ok = await create({ ...request, humanIntent: "decision", humanQuestions: questions });
|
||||
expect(ok.status).toBe(201);
|
||||
expect((await ok.json()).humanQuestions).toEqual(questions);
|
||||
const bad = await create({ ...request, humanIntent: "decision", humanQuestions: [] });
|
||||
expect(bad.status).toBe(400);
|
||||
expect(await bad.json()).toMatchObject({ error: "invalid_human_questions" });
|
||||
});
|
||||
|
||||
it("refuses questions on an update", async () => {
|
||||
await expect(repo.create({ ...request, humanIntent: "update", humanQuestions: questions })).rejects.toMatchObject({ code: "invalid_human_questions" });
|
||||
});
|
||||
|
||||
it("refuses questions sent to an agent: no human is there to click", async () => {
|
||||
await expect(repo.create({ ...request, destinationSession: "worker@rig", humanQuestions: questions })).rejects.toMatchObject({ code: "invalid_human_questions" });
|
||||
});
|
||||
|
||||
it("accepts questions when the intent is omitted: that is a decision (the CLI's documented default)", async () => {
|
||||
const created = await repo.create({ ...request, humanQuestions: questions });
|
||||
expect(repo.getById(created.qitemId)?.humanQuestions).toEqual(questions);
|
||||
});
|
||||
|
||||
it.each([
|
||||
["no questions", []],
|
||||
["five questions", Array.from({ length: 5 }, (_, i) => ({ ...questions[0]!, id: `q${i}` }))],
|
||||
["one option", [{ ...questions[0]!, options: [questions[0]!.options[0]!] }]],
|
||||
["five options", [{ ...questions[0]!, options: Array.from({ length: 5 }, (_, i) => ({ id: `o${i}`, label: `O${i}` })) }]],
|
||||
["duplicate question ids", [questions[0]!, { ...questions[1]!, id: "db" }]],
|
||||
["duplicate option ids", [{ ...questions[0]!, options: [{ id: "pg", label: "A" }, { id: "pg", label: "B" }] }]],
|
||||
["two recommended options", [{ ...questions[0]!, options: [{ id: "a", label: "A", recommended: true }, { id: "b", label: "B", recommended: true }] }]],
|
||||
["an id with spaces", [{ ...questions[0]!, id: "which db" }]],
|
||||
["an empty question", [{ ...questions[0]!, question: " " }]],
|
||||
["a label over Slack's 75-character button limit", [{ ...questions[0]!, options: [{ id: "a", label: "x".repeat(76) }, { id: "b", label: "B" }] }]],
|
||||
["not an array", { id: "db" }],
|
||||
])("refuses %s with a named error", async (_name, bad) => {
|
||||
await expect(repo.create({ ...request, humanIntent: "decision", humanQuestions: bad as never })).rejects.toMatchObject({ code: "invalid_human_questions" });
|
||||
});
|
||||
});
|
||||
|
||||
describe("recording answers", () => {
|
||||
const answer = (qitemId: string, questionId: string, optionId: string, actorSession = human) => repo.recordHumanAnswer({ qitemId, actorSession, questionId, optionId });
|
||||
|
||||
it("lets an answer change until the set is complete, then keeps the answers final for any retry", async () => {
|
||||
const { qitemId } = await repo.create({ ...request, humanIntent: "decision", humanQuestions: questions });
|
||||
expect(answer(qitemId, "db", "sqlite")).toMatchObject({ status: "recorded", complete: false });
|
||||
expect(answer(qitemId, "db", "pg")).toMatchObject({ status: "recorded", answers: { db: "pg" }, complete: false });
|
||||
expect(answer(qitemId, "ship", "yes")).toMatchObject({ status: "recorded", answers: { db: "pg", ship: "yes" }, complete: true });
|
||||
// A click after the last answer (e.g. while the hand-back is being retried) changes nothing.
|
||||
expect(answer(qitemId, "ship", "no")).toMatchObject({ status: "recorded", answers: { db: "pg", ship: "yes" }, complete: true });
|
||||
expect(repo.getById(qitemId)?.humanAnswers).toEqual({ db: "pg", ship: "yes" });
|
||||
});
|
||||
|
||||
it("does not count a question named like a built-in object property as answered", async () => {
|
||||
const builtinNamed = [{ ...questions[0]!, id: "constructor" }, questions[1]!];
|
||||
const { qitemId } = await repo.create({ ...request, humanIntent: "decision", humanQuestions: builtinNamed });
|
||||
expect(answer(qitemId, "ship", "yes")).toMatchObject({ status: "recorded", answers: { ship: "yes" }, complete: false });
|
||||
expect(answer(qitemId, "constructor", "pg")).toMatchObject({ status: "recorded", answers: { ship: "yes", constructor: "pg" }, complete: true });
|
||||
});
|
||||
|
||||
it("records nothing for another human, or once the decision is no longer pending", async () => {
|
||||
const { qitemId } = await repo.create({ ...request, humanIntent: "decision", humanQuestions: questions });
|
||||
expect(answer(qitemId, "db", "pg", "someone-else@external")).toMatchObject({ status: "not-applicable", reason: "not-the-asked-human" });
|
||||
repo.update({ qitemId, actorSession: human, state: "done", closureReason: "no-follow-on", transitionNote: "answered elsewhere" });
|
||||
expect(answer(qitemId, "db", "pg")).toMatchObject({ status: "not-applicable", reason: "state-done" });
|
||||
expect(repo.getById(qitemId)?.humanAnswers).toBeNull();
|
||||
});
|
||||
});
|
||||
|
||||
describe("rendering", () => {
|
||||
type Block = { type: string; block_id?: string; text?: { text: string }; elements?: Array<{ type: string; action_id: string; value: string; style?: string; text: { type: string; text: string } }> };
|
||||
const render = () => buildOutboundMessage({ qitemId: "q", summary: "Two quick questions", body: "Pick one for each.", humanQuestions: questions }, { sourceLabel: "proof" });
|
||||
|
||||
it("renders one button row per question, the recommended option styled primary", () => {
|
||||
const blocks = render().blocks as Block[];
|
||||
const rows = blocks.filter((b) => b.type === "actions");
|
||||
expect(rows.map((r) => r.block_id)).toEqual(["or-q:db", "or-q:ship"]);
|
||||
expect(rows[0]!.elements).toEqual([
|
||||
{ type: "button", action_id: "or-opt:pg", value: "pg", style: "primary", text: { type: "plain_text", text: "Postgres" } },
|
||||
{ type: "button", action_id: "or-opt:sqlite", value: "sqlite", text: { type: "plain_text", text: "SQLite" } },
|
||||
]);
|
||||
// Each row follows its question's text.
|
||||
const dbRow = blocks.findIndex((b) => b.block_id === "or-q:db");
|
||||
expect(blocks[dbRow - 1]?.text?.text).toContain("Which database?");
|
||||
});
|
||||
|
||||
it("keeps a complete text fallback: every question and option, the recommendation, and the typed-reply path", () => {
|
||||
const { text } = render();
|
||||
for (const s of ["Which database?", "Postgres (recommended)", "SQLite", "Ship this week?", "Yes", "No"]) expect(text).toContain(s);
|
||||
expect(text).toMatch(/reply in this thread/i);
|
||||
expect(JSON.stringify(render().blocks)).toMatch(/reply in this thread/i);
|
||||
});
|
||||
|
||||
it("escapes question and option text like any other queue-authored field", () => {
|
||||
const { blocks } = buildOutboundMessage({ qitemId: "q", humanQuestions: [{ id: "a", question: "Ping <!channel>?", options: [{ id: "x", label: "<@U1>" }, { id: "y", label: "B" }] }] }, { sourceLabel: "proof" });
|
||||
const json = JSON.stringify(blocks);
|
||||
expect(json).not.toContain("<!channel>");
|
||||
expect(json).toContain("<!channel>");
|
||||
});
|
||||
|
||||
it("renders a plain decision exactly as before", () => {
|
||||
const plain = buildOutboundMessage({ qitemId: "q", summary: "S", body: "B" }, { sourceLabel: "proof" });
|
||||
expect((plain.blocks as Block[]).some((b) => b.type === "actions")).toBe(false);
|
||||
expect(plain.text).not.toMatch(/reply in this thread/i);
|
||||
});
|
||||
});
|
||||
|
||||
describe("answering through the real Slack wire", () => {
|
||||
const reply = (body: unknown) => new Response(JSON.stringify(body), { status: 200, headers: { "content-type": "application/json" } });
|
||||
let posts: Array<Record<string, unknown>>;
|
||||
let socket: WsLike;
|
||||
let wire: ReturnType<typeof buildSlackGatewayWire>;
|
||||
beforeEach(async () => {
|
||||
const secrets = join(home, "fake.env");
|
||||
writeFileSync(secrets, "SLACK_BOT_TOKEN=xoxb-EXAMPLE-fake\nSLACK_APP_TOKEN=xapp-EXAMPLE-fake\n");
|
||||
saveConfig({ ...DEFAULT_CONFIG, enabled: true, channel: "C-TEST", secretsEnvFile: secrets, minimumLevelThatInterrupts: "NOTICE" }, home);
|
||||
posts = [];
|
||||
const sockets: WsLike[] = [];
|
||||
const contract = new MissionControlWriteContract({ db, eventBus: bus, queueRepo: repo, actionLog: new MissionControlActionLog(db) });
|
||||
const realResolve = makeHumanReplyResolver(repo, contract);
|
||||
resolveOverride = undefined;
|
||||
wire = buildSlackGatewayWire({
|
||||
home, queueRepo: repo, registry: { loadHumanRegistry: () => registry, resolveSlackHandle },
|
||||
resolveHumanReply: (input) => (resolveOverride ?? realResolve)(input),
|
||||
wsFactory: () => { const ws: WsLike = { send: () => {}, close: () => {}, onopen: null, onmessage: null, onclose: null, onerror: null }; sockets.push(ws); return ws; },
|
||||
inboundMaxConnects: 1,
|
||||
inboundRetryIntervalMs: 50, // dead-lettered clicks retry promptly
|
||||
fetchImpl: async (url, init) => {
|
||||
if (url.endsWith("apps.connections.open")) return reply({ ok: true, url: "wss://fake-slack/ws" });
|
||||
posts.push(JSON.parse(String(init?.body))); return reply({ ok: true, ts: `${posts.length}.1` });
|
||||
},
|
||||
});
|
||||
stops.push(() => wire.stop()); wire.startServices?.();
|
||||
await vi.waitFor(() => expect(sockets).toHaveLength(1));
|
||||
socket = sockets[0]!;
|
||||
socket.onopen?.();
|
||||
// Post the decision the way the outbound driver does, and wait for its receipt.
|
||||
const alert = async () => (await makeQueuePorts(repo, { loadHumanRegistry: () => registry }).listHumanAlerts({}))[0];
|
||||
decisionId = (await repo.create({ ...request, humanIntent: "decision", humanQuestions: questions })).qitemId;
|
||||
wire.dispatcher.dispatch("post_message", human, await alert());
|
||||
await vi.waitFor(async () => expect(await alert()).toBeUndefined());
|
||||
});
|
||||
let decisionId: string;
|
||||
let resolveOverride: ((input: { qitemId: string; actorSession: string; decision: string }) => Promise<"resolved" | "already-resolved" | "not-applicable">) | undefined;
|
||||
// The inbound receipt ledger: every envelope ends with one typed disposition.
|
||||
const finals = (envelopeId: string) => new InboundReceiptStore(join(home, "state", "slack-inbound-receipts.jsonl")).readAll()
|
||||
.filter((r) => r.envelopeId === envelopeId && r.status !== "received");
|
||||
const threadAcks = () => posts.filter((p) => p.thread_ts === "1.1").map((p) => String(p.text));
|
||||
|
||||
let clicks = 0;
|
||||
// A block_actions envelope, as Socket Mode delivers a button click on the root message "1.1".
|
||||
async function click(questionId: string, optionId: string, opts: { user?: string; actionTs?: string; root?: string } = {}): Promise<{ status: string; reason?: string }> {
|
||||
const actionTs = opts.actionTs ?? `${1000 + ++clicks}.5`;
|
||||
const before = finals(`e-${actionTs}`).length;
|
||||
socket.onmessage?.({ data: JSON.stringify({
|
||||
envelope_id: `e-${actionTs}`, type: "interactive",
|
||||
payload: {
|
||||
type: "block_actions", user: { id: opts.user ?? "UFOUNDER" }, channel: { id: "C-TEST" },
|
||||
container: { type: "message", message_ts: opts.root ?? "1.1", channel_id: "C-TEST" },
|
||||
message: { ts: opts.root ?? "1.1" },
|
||||
actions: [{ type: "button", block_id: `or-q:${questionId}`, action_id: `or-opt:${optionId}`, value: optionId, action_ts: actionTs }],
|
||||
},
|
||||
}) });
|
||||
await vi.waitFor(() => expect(finals(`e-${actionTs}`).length).toBe(before + 1));
|
||||
return finals(`e-${actionTs}`).at(-1)!;
|
||||
}
|
||||
const repliesToSeat = () => repo.list({ limit: 100 }).filter((q) => q.destinationSession === "author@rig");
|
||||
|
||||
it("posts the buttons, records each click, and resolves back to the seat once every question is answered", async () => {
|
||||
expect(JSON.stringify(posts[0]?.blocks)).toContain("or-opt:pg");
|
||||
|
||||
await click("db", "sqlite");
|
||||
await click("db", "pg"); // the human changes their mind before finishing
|
||||
expect(repo.getById(decisionId)).toMatchObject({ state: "pending", humanAnswers: { db: "pg" } });
|
||||
expect(repliesToSeat()).toEqual([]);
|
||||
|
||||
await click("ship", "yes");
|
||||
expect(repo.getById(decisionId)).toMatchObject({ state: "done", humanAnswers: { db: "pg", ship: "yes" } });
|
||||
const [answer] = repliesToSeat();
|
||||
expect(answer?.sourceSession).toBe(human);
|
||||
expect(answer?.tags).toEqual(expect.arrayContaining([`reply-to:${decisionId}`, "human-answer"]));
|
||||
expect(answer?.body).toContain("Which database?: Postgres");
|
||||
expect(answer?.body).toContain("Ship this week?: Yes");
|
||||
});
|
||||
|
||||
it("lands exactly one reply when the final click is replayed or clicked again", async () => {
|
||||
await click("db", "pg");
|
||||
await click("ship", "no", { actionTs: "2000.1" });
|
||||
await click("ship", "no", { actionTs: "2000.1" }); // Socket Mode redelivery
|
||||
await click("ship", "yes"); // a late click after the decision resolved
|
||||
expect(repliesToSeat()).toHaveLength(1);
|
||||
expect(repo.getById(decisionId)?.humanAnswers).toEqual({ db: "pg", ship: "no" });
|
||||
});
|
||||
|
||||
it("refuses a click from someone who is not a registered human", async () => {
|
||||
await click("db", "pg", { user: "USTRANGER" });
|
||||
expect(repo.getById(decisionId)?.humanAnswers).toBeNull();
|
||||
});
|
||||
|
||||
it("ignores a click naming a question or option the decision does not have", async () => {
|
||||
await click("db", "mysql");
|
||||
await click("region", "eu");
|
||||
expect(repo.getById(decisionId)?.humanAnswers).toBeNull();
|
||||
});
|
||||
|
||||
it("with supplemental detail, the buttons ride the root post only, and a click there answers the decision", async () => {
|
||||
const withDetail = await repo.create({ ...request, humanIntent: "decision", humanQuestions: [questions[0]!], humanDetail: "Benchmarks: Postgres 2x faster on our load." });
|
||||
const alert = (await makeQueuePorts(repo, { loadHumanRegistry: () => registry }).listHumanAlerts({})).find((q) => q.qitemId === withDetail.qitemId);
|
||||
wire.dispatcher.dispatch("post_message", human, alert);
|
||||
await vi.waitFor(() => expect(posts).toHaveLength(3));
|
||||
const [root, detail] = [posts[1]!, posts[2]!];
|
||||
expect(root.thread_ts).toBeUndefined();
|
||||
expect(JSON.stringify(root.blocks)).toContain("or-opt:pg");
|
||||
expect(detail.thread_ts).toBe("2.1");
|
||||
expect(JSON.stringify(detail.blocks)).not.toContain("or-opt:");
|
||||
expect(String(detail.text)).not.toContain("Which database?");
|
||||
|
||||
await click("db", "pg", { root: "2.1" });
|
||||
expect(repo.getById(withDetail.qitemId)).toMatchObject({ state: "done", humanAnswers: { db: "pg" } });
|
||||
expect(repo.getById(decisionId)?.humanAnswers).toBeNull();
|
||||
});
|
||||
|
||||
it("confirms each click in the decision's thread, naming what is still unanswered", async () => {
|
||||
await click("db", "pg");
|
||||
await vi.waitFor(() => expect(threadAcks()).toHaveLength(1));
|
||||
expect(threadAcks()[0]).toBe("Recorded: Which database?: Postgres. Still to answer: Ship this week?");
|
||||
await click("ship", "yes");
|
||||
await vi.waitFor(() => expect(threadAcks()).toHaveLength(2));
|
||||
expect(threadAcks()[1]).toBe("All answered, sent back: Which database?: Postgres; Ship this week?: Yes");
|
||||
});
|
||||
|
||||
it("redacts secret-like question and option text in the click confirmations, as the question post does", async () => {
|
||||
const secretQuestions = [
|
||||
{ id: "tok", question: "Rotate xoxb-LEAK-question?", options: [{ id: "a", label: "Use xoxb-LEAK-option" }, { id: "b", label: "Skip" }] },
|
||||
{ id: "ship", question: "Ship this week?", options: [{ id: "yes", label: "Yes" }, { id: "no", label: "No" }] },
|
||||
];
|
||||
const secret = await repo.create({ ...request, humanIntent: "decision", humanQuestions: secretQuestions });
|
||||
const alert = async () => (await makeQueuePorts(repo, { loadHumanRegistry: () => registry }).listHumanAlerts({})).find((q) => q.qitemId === secret.qitemId);
|
||||
wire.dispatcher.dispatch("post_message", human, await alert());
|
||||
await vi.waitFor(async () => expect(await alert()).toBeUndefined()); // posted and its thread mapped
|
||||
const acks = () => posts.filter((p) => p.thread_ts === "2.1").map((p) => String(p.text));
|
||||
|
||||
await click("tok", "a", { root: "2.1" });
|
||||
await click("ship", "yes", { root: "2.1" });
|
||||
await vi.waitFor(() => expect(acks()).toHaveLength(2));
|
||||
expect(acks().join("\n")).not.toContain("xoxb-LEAK");
|
||||
expect(acks()[1]).toContain("[redacted-secret]");
|
||||
});
|
||||
|
||||
it("recovers a failed hand-back without another click: the dead-letter retry lands it", async () => {
|
||||
const create = repo.create.bind(repo);
|
||||
vi.spyOn(repo, "create").mockImplementationOnce(async () => { throw new Error("database is locked"); }).mockImplementation(create);
|
||||
await click("db", "pg");
|
||||
expect(await click("ship", "yes")).toMatchObject({ status: "handler-failed" });
|
||||
await vi.waitFor(() => expect(repo.getById(decisionId)?.state).toBe("done"));
|
||||
expect(repliesToSeat()).toHaveLength(1);
|
||||
});
|
||||
|
||||
it("does not report success when the resolve does not apply", async () => {
|
||||
resolveOverride = async () => "not-applicable";
|
||||
await click("db", "pg");
|
||||
expect(await click("ship", "yes")).toMatchObject({ status: "refused", reason: "resolve-not-applicable" });
|
||||
expect(threadAcks().some((t) => t.startsWith("All answered"))).toBe(false);
|
||||
});
|
||||
|
||||
it("still resolves on a typed reply in the thread (the \"Other\" answer)", async () => {
|
||||
socket.onmessage?.({ data: JSON.stringify({ envelope_id: "e-typed", type: "events_api", payload: { event: { type: "message", user: "UFOUNDER", text: "Neither — use DuckDB", ts: "3000.1", thread_ts: "1.1", channel: "C-TEST" } } }) });
|
||||
await vi.waitFor(() => expect(repo.getById(decisionId)?.state).toBe("done"));
|
||||
expect(repliesToSeat()[0]?.body).toContain("Neither — use DuckDB");
|
||||
});
|
||||
});
|
||||
});
|
||||
@@ -6,7 +6,7 @@
|
||||
"core/agent-startup-and-context-ingestion/SKILL.md": "ed9a3fba75d9c53d44d3fb2a9b4703a1bdb5d13e78e4e34f18c2887c23c9d248",
|
||||
"core/cross-host-rig-commands/SKILL.md": "844a578d0f51cf8f02f8bfd7ffb3bbd6a37d511f721266c0d758c49157c1447a",
|
||||
"core/human-in-the-loop/SKILL.md": "14cf0c4a0dc89179121465d4b8d81d40406427351adcdc34f875097d1e6e9ab2",
|
||||
"core/messaging-the-human/SKILL.md": "4ddd7d4dbcf4f9d5f7beb01d9a98db1766ebc7d52e5636786f144f14f9345b9f",
|
||||
"core/messaging-the-human/SKILL.md": "581372434f27b1994d516e144be1324c7c5981cacf64fef6f09d1a57017faa08",
|
||||
"core/openrig-architect/SKILL.md": "0c3f1b7ff5bf185c9da963905229e951ffe1fb6d45a3cc9247cb747fc1d9a59d",
|
||||
"core/openrig-cmux/SKILL.md": "73a86eac45d68b16a38fc0a4438713c634d86061073dac96a4d9e656ea5d34dc",
|
||||
"core/openrig-herdr/SKILL.md": "37508a02f50bb69ee97c22c37ccf70aa8e37dc0784061eb2a5ad813571eebd1a",
|
||||
@@ -67,7 +67,7 @@
|
||||
"forming-an-openrig-mental-model/SKILL.md": "cff1ca1640fec65b180539a2c6c3c9044d73aa02f9d13d0d24763446ac974da0",
|
||||
"loading-addressable-markdown/SKILL.md": "104dc6a02ae6e2d0eff519387c6839cb09803945b9ee35ec3498af8feaaa8d1e",
|
||||
"loading-addressable-markdown/scripts/resolve-markdown.mjs": "ef113af7bb297058a8c23a82a4c171232ae8c2d6b1705a89a80140000ad3bca0",
|
||||
"messaging-the-human/SKILL.md": "4ddd7d4dbcf4f9d5f7beb01d9a98db1766ebc7d52e5636786f144f14f9345b9f",
|
||||
"messaging-the-human/SKILL.md": "581372434f27b1994d516e144be1324c7c5981cacf64fef6f09d1a57017faa08",
|
||||
"mission-slice-sop/SKILL.md": "3c3783d2f2d05cead97c94d10a9af00b397b810a8fdb7d599930fc1ea226305d",
|
||||
"openrig-operating-model/SKILL.md": "53ec4d9c7764f5bc23616885d193b79ddb1be139fd010f4b89f143d39dfdcaa6",
|
||||
"openrig-operating-model/scripts/compose.py": "332a13580393fd1421f0b498eae51f01243de31ba168a86db52a6251fc4cc384",
|
||||
@@ -102,7 +102,7 @@
|
||||
"core/agent-startup-and-context-ingestion/SKILL.md": "ed9a3fba75d9c53d44d3fb2a9b4703a1bdb5d13e78e4e34f18c2887c23c9d248",
|
||||
"core/cross-host-rig-commands/SKILL.md": "844a578d0f51cf8f02f8bfd7ffb3bbd6a37d511f721266c0d758c49157c1447a",
|
||||
"core/human-in-the-loop/SKILL.md": "14cf0c4a0dc89179121465d4b8d81d40406427351adcdc34f875097d1e6e9ab2",
|
||||
"core/messaging-the-human/SKILL.md": "4ddd7d4dbcf4f9d5f7beb01d9a98db1766ebc7d52e5636786f144f14f9345b9f",
|
||||
"core/messaging-the-human/SKILL.md": "581372434f27b1994d516e144be1324c7c5981cacf64fef6f09d1a57017faa08",
|
||||
"core/openrig-architect/SKILL.md": "0c3f1b7ff5bf185c9da963905229e951ffe1fb6d45a3cc9247cb747fc1d9a59d",
|
||||
"core/openrig-cmux/SKILL.md": "73a86eac45d68b16a38fc0a4438713c634d86061073dac96a4d9e656ea5d34dc",
|
||||
"core/openrig-herdr/SKILL.md": "37508a02f50bb69ee97c22c37ccf70aa8e37dc0784061eb2a5ad813571eebd1a",
|
||||
|
||||
@@ -92,6 +92,17 @@ the human (a pending human decision, or a row parked on the human), since a
|
||||
reply in that thread would answer the decision. In every such case the
|
||||
`--verify` result says `threaded: false` with the reason.
|
||||
|
||||
A **decision** with a few clear choices can carry `--human-questions-file <path>`:
|
||||
a JSON array of 1–4 questions, each
|
||||
`{"id", "question", "options": [{"id", "label", "recommended"?}]}` with 2–4
|
||||
options (labels up to 75 characters, at most one recommended). Slack shows each
|
||||
question as a row of buttons. Each click records that answer on the item, and
|
||||
the decision resolves once every question has one. You then receive one reply
|
||||
row listing the answers, and the item's `humanAnswers` holds the option ids. The
|
||||
human may instead type a reply in the thread; that resolves the decision as
|
||||
usual, so read the reply rather than assuming an option was picked. Keep the
|
||||
brief complete: the questions add buttons, they do not replace the explanation.
|
||||
|
||||
If an existing agent-owned row must wait for a **decision**, block it on the **new live qitem ID**
|
||||
(`rig queue block <work-id> --on <human-qitem-id> ...`), not on the human address.
|
||||
Completion of the human qitem resumes its dependants. Blocking on the human as
|
||||
|
||||
Reference in New Issue
Block a user