diff --git a/docs/superpowers/specs/2026-08-04-workflow-email-connector-demo-design.md b/docs/superpowers/specs/2026-08-04-workflow-email-connector-demo-design.md new file mode 100644 index 0000000..2b24b38 --- /dev/null +++ b/docs/superpowers/specs/2026-08-04-workflow-email-connector-demo-design.md @@ -0,0 +1,41 @@ +# Maestro Workflow 이메일 커넥터 데모 설계 (목 드라이버) + +- 날짜: 2026-08-04 +- 상태: 확정 (사용자: 목 드라이버 데모 먼저 — IMAP/Gmail은 드라이버 교체로 후속) +- 범위: `workflow/examples/` + `workflow/tests/` (공개 계약만 사용하는 참조 클라이언트) + +## 0. 목표 + +비전 §4(d)의 이메일 업무 루프 전체를 실서버로 증명한다: +받은편지함 → email-triage 요청(체인 루트) → 운영자 승인 → email-reply +요청(체인 연결, 초안 payload) → 운영자 승인 → **발송 실행(드라이버)** → ack. +Workflow 서버는 무변경 — 커넥터는 기존 계약(actor 등록·요청·WS 구독·ack)만 쓴다. + +## 1. 구조 + +- `workflow/examples/email-connector/lib.mjs`: + `runEmailConnector({ serverUrl, serverToken, actorId, driver, log, decisionTimeoutMs })` + → 처리 요약 `{ chains, sent, skipped }` 반환. 드라이버 인터페이스: + `listUnprocessed()` / `send({ to, subject, body })` / `markProcessed(id)`. +- `mockInbox.mjs`: 가짜 메일 2통 + 발송 기록(`sent[]`)을 가진 목 드라이버. +- `connector.mjs`: CLI 래퍼 (`node connector.mjs` — env로 서버/토큰 지정). +- 결정 수신: actor 토큰 WS 구독(WORKFLOW_AUTH → 자기 WORKFLOW_DECIDED)을 + 1차로, 폴링(GET .../decision)을 보조로 사용 — 스펙 2026-08-04(actor WS)의 + "WS는 알림, 폴링이 보장" 원칙 그대로. +- 반려(reject/revise) 시: 해당 메일은 발송하지 않고 `skipped`로 기록(체인에 + 결정 사유가 남는다). 데모 범위에서 재초안 루프는 생략. + +## 2. e2e 테스트 (`workflow/tests/email-connector-demo.test.mjs`) + +실서버(엄격 모드)를 띄우고 커넥터를 병행 실행, 테스트가 운영자 역할로 +pending 요청을 폴링·승인한다. 단언: + +- 메일 2통 → 체인 2개, 각 체인 = [email-triage, email-reply] (chain API 검증) +- 승인 완료 후 목 드라이버 `sent`에 답장 2건, 수신자/제목 일치 +- 결정 2×2건 모두 ack 상태(delivery.status)로 종결 +- 반려 시나리오 1건: triage 반려 → reply 미생성·발송 0건·skipped 기록 + +## 3. 비범위 + +실제 IMAP/Gmail 드라이버(후속 — 인터페이스 동일), 재초안 반복 루프, +대시보드 체인 시각화. diff --git a/workflow/README.md b/workflow/README.md index b39f3f5..318a28c 100644 --- a/workflow/README.md +++ b/workflow/README.md @@ -40,6 +40,11 @@ actor도 자신의 actorToken으로 같은 `WORKFLOW_AUTH` 핸드셰이크를 - 스펙: [`docs/superpowers/specs/2026-08-03-workflow-strict-dashboard-design.md`](../docs/superpowers/specs/2026-08-03-workflow-strict-dashboard-design.md) +## 예제 + +- [`examples/email-connector/`](examples/email-connector/) — 이메일 업무 루프 + 참조 클라이언트(목 드라이버). e2e: `tests/email-connector-demo.test.mjs`. + ## 알려진 한계 (MVP) - 토큰은 localStorage에 평문 저장된다 — 로컬 신뢰 기기 전제. TLS 없음, 기본 diff --git a/workflow/examples/email-connector/README.md b/workflow/examples/email-connector/README.md new file mode 100644 index 0000000..1148ae2 --- /dev/null +++ b/workflow/examples/email-connector/README.md @@ -0,0 +1,21 @@ +# 이메일 커넥터 (참조 클라이언트, 목 드라이버) + +비전 §4(d)의 이메일 업무 루프 데모: 받은편지함 → `email-triage`(체인 루트) +→ 승인 → `email-reply`(체인) → 승인 → 발송(드라이버) → ack. +Workflow 공개 계약만 사용하며, 서버 코드는 무변경이다. + +## 실행 + + npm run server # 터미널 1 (workflow/) + node examples/email-connector/connector.mjs # 터미널 2 + +엄격 모드면 두 터미널 모두 `MAESTRO_WORKFLOW_SERVER_TOKEN`을 지정한다. +대시보드(레인)에서 ✉/↩ 프리셋 카드를 승인·반려하면 커넥터가 WS로 결정을 +받아 발송을 시뮬레이션하고 요약을 출력한다. + +## 실제 이메일 연결 (후속) + +`mockInbox.mjs`와 같은 인터페이스(`listUnprocessed` / `send` / +`markProcessed`)로 IMAP 또는 Gmail API 드라이버를 구현해 `connector.mjs`에서 +교체하면 된다 — 커넥터 본체(`lib.mjs`)는 수정 불필요. 이메일 자격증명은 +커넥터 프로세스에만 머물고 Workflow 서버에는 절대 전달되지 않는다. diff --git a/workflow/examples/email-connector/connector.mjs b/workflow/examples/email-connector/connector.mjs new file mode 100644 index 0000000..126c3ff --- /dev/null +++ b/workflow/examples/email-connector/connector.mjs @@ -0,0 +1,20 @@ +#!/usr/bin/env node +// 이메일 커넥터 CLI (목 드라이버). 실행 전 workflow 서버가 떠 있어야 한다. +// MAESTRO_WORKFLOW_SERVER_TOKEN=<서버토큰> node connector.mjs +// 서버 토큰은 actor 등록에만 쓰이고, 이후는 발급받은 actor 토큰으로만 통신한다. +import { runEmailConnector } from './lib.mjs'; +import { createMockInboxDriver } from './mockInbox.mjs'; + +const serverUrl = process.env.MAESTRO_WORKFLOW_SERVER_URL || 'http://127.0.0.1:8090'; +const serverToken = process.env.MAESTRO_WORKFLOW_SERVER_TOKEN || ''; + +const driver = createMockInboxDriver(); +const summary = await runEmailConnector({ + serverUrl, + serverToken, + driver, + log: console.log, +}); + +console.log('\n── 처리 요약 ──'); +console.log(JSON.stringify(summary, null, 2)); diff --git a/workflow/examples/email-connector/lib.mjs b/workflow/examples/email-connector/lib.mjs new file mode 100644 index 0000000..4605909 --- /dev/null +++ b/workflow/examples/email-connector/lib.mjs @@ -0,0 +1,164 @@ +// 이메일 커넥터 참조 클라이언트 (스펙 2026-08-04 데모 §1). +// Workflow 공개 계약만 사용한다: actor 등록 → 요청 생성(체인) → WS 구독(1차)/폴링(보조) → ack. +// 발송 실행은 드라이버 몫 — Workflow는 record-only 그대로다. +import WebSocket from 'ws'; + +export async function runEmailConnector({ + serverUrl, + serverToken, + actorId = 'agent_email', + driver, + log = () => {}, + decisionTimeoutMs = 60000, +}) { + if (!driver) { + throw new Error('driver가 필요합니다 (listUnprocessed/send/markProcessed)'); + } + + const httpJson = async (path, { method = 'GET', token = '', body = null } = {}) => { + const res = await fetch(`${serverUrl}${path}`, { + method, + headers: { + 'Content-Type': 'application/json', + ...(token ? { Authorization: `Bearer ${token}` } : {}), + }, + body: body ? JSON.stringify(body) : undefined, + }); + const parsed = await res.json().catch(() => ({})); + if (!res.ok) { + throw new Error(parsed.error || `HTTP ${res.status} ${path}`); + } + return parsed; + }; + + // 1. actor 등록 (서버 토큰) → actor 토큰 확보 + const registration = await httpJson('/api/actors/register', { + method: 'POST', + token: serverToken, + body: { actorId }, + }); + const actorToken = registration.actorToken; + + // 2. WS 구독 — 자기 결정(WORKFLOW_DECIDED)만 수신 (actor 스코프) + const decisionWaiters = new Map(); // requestId → resolve + const receivedDecisions = new Map(); // requestId → { item, request } + const ws = new WebSocket(serverUrl.replace(/^http/, 'ws')); + await new Promise((resolve, reject) => { + ws.on('open', () => { + ws.send(JSON.stringify({ type: 'WORKFLOW_AUTH', token: actorToken })); + }); + ws.on('message', (raw) => { + const data = JSON.parse(raw.toString()); + if (data.type === 'WORKFLOW_AUTH_OK') { + resolve(); + return; + } + if (data.type === 'WORKFLOW_DECIDED' && data.request) { + receivedDecisions.set(data.request.requestId, { item: data.item, request: data.request }); + const waiter = decisionWaiters.get(data.request.requestId); + if (waiter) waiter({ item: data.item, request: data.request }); + } + }); + ws.on('error', reject); + }); + + // WS가 1차, 폴링이 보장 (재연결·놓침 복구용 보조 경로) + const awaitDecision = async (requestId) => { + if (receivedDecisions.has(requestId)) { + return receivedDecisions.get(requestId); + } + return new Promise((resolve, reject) => { + const deadline = setTimeout(() => { + clearInterval(poller); + decisionWaiters.delete(requestId); + reject(new Error(`결정 대기 시간 초과: ${requestId}`)); + }, decisionTimeoutMs); + const settle = (decision) => { + clearTimeout(deadline); + clearInterval(poller); + decisionWaiters.delete(requestId); + resolve(decision); + }; + decisionWaiters.set(requestId, settle); + const poller = setInterval(async () => { + try { + const status = await httpJson(`/api/decision-requests/${encodeURIComponent(requestId)}/decision`, { token: actorToken }); + if (status.item) { + settle({ item: status.item, request: { requestId } }); + } + } catch { + // 폴링 실패는 다음 주기에 재시도 + } + }, 1500); + }); + }; + + const createRequest = (payload) => httpJson('/api/decision-requests', { + method: 'POST', + token: actorToken, + body: payload, + }); + + const ackDecision = (decisionId) => httpJson(`/api/decisions/${encodeURIComponent(decisionId)}/ack`, { + method: 'POST', + token: actorToken, + }); + + // 3. 메일별 체인 처리: triage(루트) → 승인 시 reply(체인) → 승인 시 발송 + const summary = { chains: [], sent: [], skipped: [] }; + + try { + for (const mail of driver.listUnprocessed()) { + log(`📥 처리 시작: ${mail.subject} (${mail.from})`); + const triage = await createRequest({ + subjectType: 'email-triage', + subject: { + title: `메일 분류: ${mail.subject}`, + summary: mail.body, + payload: { from: mail.from, subject: mail.subject, proposedAction: mail.proposedAction }, + }, + }); + const triageDecision = await awaitDecision(triage.item.requestId); + await ackDecision(triageDecision.item.decisionId); + + if (triageDecision.item.decision !== 'approve') { + summary.skipped.push({ mailId: mail.id, stage: 'triage', decision: triageDecision.item.decision }); + driver.markProcessed(mail.id); + log(`⏭ 분류 단계에서 중단(${triageDecision.item.decision}): ${mail.subject}`); + continue; + } + + const reply = await createRequest({ + subjectType: 'email-reply', + parentRequestId: triage.item.requestId, + subject: { + title: `답장 승인: ${mail.subject}`, + summary: mail.draftReply, + payload: { to: mail.from, subject: `Re: ${mail.subject}`, draft: mail.draftReply }, + }, + }); + const replyDecision = await awaitDecision(reply.item.requestId); + await ackDecision(replyDecision.item.decisionId); + + if (replyDecision.item.decision === 'approve') { + driver.send({ to: mail.from, subject: `Re: ${mail.subject}`, body: mail.draftReply }); + summary.sent.push({ mailId: mail.id, to: mail.from }); + log(`📤 발송 완료: Re: ${mail.subject} → ${mail.from}`); + } else { + summary.skipped.push({ mailId: mail.id, stage: 'reply', decision: replyDecision.item.decision }); + log(`⏭ 답장 반려(${replyDecision.item.decision}): ${mail.subject}`); + } + + driver.markProcessed(mail.id); + summary.chains.push({ + mailId: mail.id, + triageRequestId: triage.item.requestId, + replyRequestId: reply.item.requestId, + }); + } + } finally { + ws.close(); + } + + return summary; +} diff --git a/workflow/examples/email-connector/mockInbox.mjs b/workflow/examples/email-connector/mockInbox.mjs new file mode 100644 index 0000000..031e181 --- /dev/null +++ b/workflow/examples/email-connector/mockInbox.mjs @@ -0,0 +1,38 @@ +// 목 받은편지함 드라이버 (스펙 2026-08-04 데모 §1): 자격증명 없이 전체 루프를 증명한다. +// 실제 IMAP/Gmail 드라이버는 이 인터페이스(listUnprocessed/send/markProcessed)만 맞추면 된다. + +export function createMockInboxDriver(seedMails = null) { + const mails = seedMails || [ + { + id: 'mail-1', + from: 'client@corp.com', + subject: '견적 회신 요청', + body: '지난주 논의한 범위로 견적 부탁드립니다.', + proposedAction: '표준 견적 템플릿으로 회신', + draftReply: '안녕하세요, 요청하신 견적을 첨부와 같이 회신드립니다. 감사합니다.', + }, + { + id: 'mail-2', + from: 'partner@vendor.io', + subject: '미팅 일정 조율', + body: '다음 주 수요일 오후 가능하신가요?', + proposedAction: '수요일 15시로 수락 회신', + draftReply: '안녕하세요, 수요일 15시 좋습니다. 초대장 보내주시면 참석하겠습니다.', + }, + ]; + const processed = new Set(); + const sent = []; + + return { + sent, + listUnprocessed() { + return mails.filter((mail) => !processed.has(mail.id)); + }, + send({ to, subject, body }) { + sent.push({ to, subject, body }); + }, + markProcessed(mailId) { + processed.add(mailId); + }, + }; +} diff --git a/workflow/tests/email-connector-demo.test.mjs b/workflow/tests/email-connector-demo.test.mjs new file mode 100644 index 0000000..666cf00 --- /dev/null +++ b/workflow/tests/email-connector-demo.test.mjs @@ -0,0 +1,123 @@ +// 이메일 커넥터 e2e (스펙 2026-08-04 데모 §2): 실서버 + 커넥터 + 스크립트 운영자로 전체 루프 검증. +import test from 'node:test'; +import assert from 'node:assert/strict'; +import { startServer, cleanupDataDir, authHeaders } from './helpers.mjs'; +import { runEmailConnector } from '../examples/email-connector/lib.mjs'; +import { createMockInboxDriver } from '../examples/email-connector/mockInbox.mjs'; + +const SERVER_TOKEN = 'wf-server-secret'; + +// 운영자 역할: pending 요청을 폴링해 decide 함수로 결정한다. done()이 true가 되면 즉시 종료. +async function operateUntil(server, decide, { timeoutMs = 20000, done = () => false } = {}) { + const deadline = Date.now() + timeoutMs; + const decided = new Set(); + while (Date.now() < deadline && !done()) { + const res = await fetch(`http://127.0.0.1:${server.port}/api/decision-requests?status=pending_decision`, { + headers: authHeaders(SERVER_TOKEN), + }); + const { items } = await res.json(); + for (const request of items) { + if (decided.has(request.requestId)) continue; + const decision = decide(request); + if (!decision) continue; + await fetch(`http://127.0.0.1:${server.port}/api/decision-requests/${request.requestId}/decide`, { + method: 'POST', + headers: { 'Content-Type': 'application/json', ...authHeaders(SERVER_TOKEN) }, + body: JSON.stringify(decision), + }); + decided.add(request.requestId); + } + await new Promise((resolve) => setTimeout(resolve, 150)); + } +} + +test('전체 루프: 메일 2통이 triage→reply 체인으로 승인되어 발송·ack까지 완결된다', async () => { + const server = await startServer({ serverToken: SERVER_TOKEN }); + const driver = createMockInboxDriver(); + try { + let connectorDone = false; + const operator = operateUntil(server, () => ({ decision: 'approve', comment: '' }), { + timeoutMs: 15000, + done: () => connectorDone, + }); + const summary = await runEmailConnector({ + serverUrl: `http://127.0.0.1:${server.port}`, + serverToken: SERVER_TOKEN, + driver, + decisionTimeoutMs: 15000, + }); + connectorDone = true; + + assert.equal(summary.chains.length, 2); + assert.equal(summary.sent.length, 2); + assert.equal(summary.skipped.length, 0); + assert.deepEqual(driver.sent.map((mail) => mail.to), ['client@corp.com', 'partner@vendor.io']); + assert.ok(driver.sent[0].subject.startsWith('Re: ')); + + // 체인 API로 triage→reply 연결 검증 + for (const chain of summary.chains) { + const chainRes = await fetch( + `http://127.0.0.1:${server.port}/api/decision-requests/${chain.replyRequestId}/chain`, + { headers: authHeaders(SERVER_TOKEN) }, + ); + const { items } = await chainRes.json(); + assert.deepEqual( + items.map((item) => item.requestId), + [chain.triageRequestId, chain.replyRequestId], + ); + assert.deepEqual(items.map((item) => item.subjectType), ['email-triage', 'email-reply']); + } + + // 결정 4건 모두 ack로 종결됐는지 (actor 폴링 관점) + const historyRes = await fetch(`http://127.0.0.1:${server.port}/api/history?limit=60`, { + headers: authHeaders(SERVER_TOKEN), + }); + const historyItems = (await historyRes.json()).items; + const ackCount = historyItems.filter((entry) => entry.event === 'ACKNOWLEDGED').length; + assert.equal(ackCount, 4); + + await operator; + } finally { + await server.stop(); + cleanupDataDir(server.dataDir); + } +}); + +test('반려 루프: triage 반려 시 reply를 만들지 않고 발송 0건으로 기록한다', async () => { + const server = await startServer({ serverToken: SERVER_TOKEN }); + const driver = createMockInboxDriver([ + { + id: 'mail-spam', + from: 'spam@junk.io', + subject: '광고: 무제한 크레딧', + body: '지금 바로 구매하세요!', + proposedAction: '무시 또는 스팸 처리', + draftReply: '(초안 없음)', + }, + ]); + try { + let connectorDone = false; + const operator = operateUntil(server, () => ({ decision: 'reject', comment: '스팸 — 회신 불필요' }), { + timeoutMs: 10000, + done: () => connectorDone, + }); + const summary = await runEmailConnector({ + serverUrl: `http://127.0.0.1:${server.port}`, + serverToken: SERVER_TOKEN, + driver, + decisionTimeoutMs: 10000, + }); + connectorDone = true; + + assert.equal(summary.sent.length, 0); + assert.equal(summary.chains.length, 0); + assert.deepEqual(summary.skipped, [{ mailId: 'mail-spam', stage: 'triage', decision: 'reject' }]); + assert.equal(driver.sent.length, 0); + assert.equal(driver.listUnprocessed().length, 0); // 반려도 처리 완료로 표시 + + await operator; + } finally { + await server.stop(); + cleanupDataDir(server.dataDir); + } +});