diff --git a/sdk/src/index.ts b/sdk/src/index.ts index 59875f542..029fb87e6 100644 --- a/sdk/src/index.ts +++ b/sdk/src/index.ts @@ -144,3 +144,13 @@ export { type Fetcher, type PollOptions, } from './hn-poller.js'; + +// Linear adapter — second proactive workload on gate 2 primitives, same +// pattern as hn-poller (adapter outside kernel/, submits via journal). +export { + pollLinearOnce, + LINEAR_GRAPHQL_ENDPOINT, + type LinearFetcher, + type LinearIssueRef, + type PollOptions as LinearPollOptions, +} from './linear-poller.js'; diff --git a/sdk/src/linear-poller.ts b/sdk/src/linear-poller.ts new file mode 100644 index 000000000..b86e7a546 --- /dev/null +++ b/sdk/src/linear-poller.ts @@ -0,0 +1,149 @@ +/** + * Linear -> relayflow events. + * + * Same shape as sdk/src/hn-poller.ts: this adapter runs on the authoring + * surface (not kernel) and submits events through the journal protocol + * (`event.submit`). The kernel learns about Linear the way it learns about + * everything else: as an event. This preserves PR #16's settled decision + * that provider-specific product logic + network I/O stay out of `kernel/`. + * + * Second proactive workload on gate 2 primitives (hn-monitor is the first), + * per RFC-0001 §3 gate 2 ("hn-monitor or linear"). Proves the pattern + * generalizes beyond the HN adapter. + * + * Auth: Linear requires an API token, unlike HN's public feed. The token is + * read from `LINEAR_API_TOKEN` (env-configurable via the runner or CLI). In + * a gate-6 world this would come from a relayfile-mounted credential, but + * for gate-2 scope we accept the env-var bootstrap and note it in the flow's + * preflight. + * + * Deduplication is the kernel's job (flow's `dedupeKeyTemplate` + the + * (flow, subscription, key) claim). This function submits every issue it + * fetches without tracking what it's seen. + */ + +const LINEAR_GRAPHQL_URL = 'https://api.linear.app/graphql'; +const DEFAULT_ISSUE_LIMIT = 10; + +/** Anything that can submit an event through the journal protocol. */ +export interface EventSink { + eventSubmit(spec: unknown, event: { type: string; payload?: unknown; key?: string }): Promise; +} + +/** Injected so parsing and submission stay deterministic in tests. */ +export type LinearFetcher = (query: string, token: string) => Promise; + +const defaultFetcher: LinearFetcher = async (query, token) => { + const response = await fetch(LINEAR_GRAPHQL_URL, { + method: 'POST', + headers: { + 'Content-Type': 'application/json', + Authorization: token, + }, + body: JSON.stringify({ query }), + }); + if (!response.ok) { + throw new Error(`Linear fetch failed: HTTP ${response.status}`); + } + return response.text(); +}; + +/** Minimal issue shape used to construct the wake event payload. */ +export interface LinearIssueRef { + id: string; + identifier?: string; + title?: string; + createdAt?: string; + url?: string; +} + +export interface PollOptions { + /** Cap on issues per poll. Default 10. */ + issueLimit?: number; + /** Fetcher override (tests). */ + fetcher?: LinearFetcher; + /** API token override (production reads LINEAR_API_TOKEN env var). */ + token?: string; + /** ISO timestamp filter — only fetch issues created after this. */ + createdAfter?: string; +} + +/** + * Fetch the newest Linear issues once and submit each as an event. + * + * Journal errors propagate; fetch errors propagate (the runner layer + * decides what to swallow). Empty result is not an error. + */ +export async function pollLinearOnce( + spec: unknown, + sink: EventSink, + options: PollOptions = {}, +): Promise { + const issueLimit = options.issueLimit ?? DEFAULT_ISSUE_LIMIT; + const fetcher = options.fetcher ?? defaultFetcher; + const token = options.token ?? process.env['LINEAR_API_TOKEN'] ?? ''; + if (token === '') { + throw new Error( + 'linear-poller: LINEAR_API_TOKEN is not set — provide options.token or set the env var', + ); + } + + const createdAfter = options.createdAfter; + // GraphQL: fetch newest issues, optionally filtered by createdAt. + const filter = createdAfter + ? `, filter: { createdAt: { gt: "${createdAfter}" } }` + : ''; + const query = ` + query { + issues(first: ${issueLimit}, orderBy: createdAt${filter}) { + nodes { + id + identifier + title + createdAt + url + } + } + } + `.trim(); + + const body = await fetcher(query, token); + + let parsed: unknown; + try { + parsed = JSON.parse(body); + } catch (cause) { + throw new Error(`Linear response was not JSON: ${String(cause)}`); + } + + // Shape check — Linear returns { data: { issues: { nodes: [...] } } }; + // GraphQL errors surface as { errors: [...] }. + if (parsed && typeof parsed === 'object' && 'errors' in parsed) { + throw new Error(`Linear GraphQL error: ${JSON.stringify((parsed as any).errors)}`); + } + const nodes = (parsed as any)?.data?.issues?.nodes; + if (!Array.isArray(nodes)) { + throw new Error('Linear response was missing data.issues.nodes array'); + } + + const outcomes: unknown[] = []; + for (const node of nodes as LinearIssueRef[]) { + if (!node || typeof node.id !== 'string') continue; + outcomes.push( + await sink.eventSubmit(spec, { + type: 'linear.issue_created', + payload: { + id: node.id, + identifier: node.identifier, + title: node.title, + created_at: node.createdAt, + url: node.url, + type: 'issue', + }, + }), + ); + } + return outcomes; +} + +export const LINEAR_GRAPHQL_ENDPOINT = LINEAR_GRAPHQL_URL; diff --git a/sdk/tests/linear-poller.test.ts b/sdk/tests/linear-poller.test.ts new file mode 100644 index 000000000..df8bcb347 --- /dev/null +++ b/sdk/tests/linear-poller.test.ts @@ -0,0 +1,102 @@ +import { describe, expect, it } from 'vitest'; +import { pollLinearOnce, LINEAR_GRAPHQL_ENDPOINT } from '../src/linear-poller.js'; + +/** Recorded response — the test never touches the network. */ +const RECORDED_ISSUES = JSON.stringify({ + data: { + issues: { + nodes: [ + { id: 'issue-uuid-1', identifier: 'ENG-1234', title: 'Fix login flow', createdAt: '2026-08-31T10:00:00Z', url: 'https://linear.app/x/issue/ENG-1234' }, + { id: 'issue-uuid-2', identifier: 'ENG-1235', title: 'Add dashboard', createdAt: '2026-08-31T10:05:00Z', url: 'https://linear.app/x/issue/ENG-1235' }, + ], + }, + }, +}); + +function recordingSink() { + const submitted: Array<{ spec: unknown; event: { type: string; payload?: unknown } }> = []; + return { + submitted, + async eventSubmit(spec: unknown, event: { type: string; payload?: unknown }) { + submitted.push({ spec, event }); + return { matched: true, deduped: false }; + }, + }; +} + +describe('linear poller', () => { + it('submits one event per issue, up to the limit, through the journal protocol', async () => { + const sink = recordingSink(); + const spec = { name: 'linear-monitor' }; + + await pollLinearOnce(spec, sink, { + token: 'lin_api_test', + fetcher: async (query, token) => { + expect(token).toBe('lin_api_test'); + expect(query).toContain('issues('); + expect(query).toContain('orderBy: createdAt'); + return RECORDED_ISSUES; + }, + }); + + expect(sink.submitted).toHaveLength(2); + expect(sink.submitted[0].event.type).toBe('linear.issue_created'); + expect((sink.submitted[0].event.payload as any).id).toBe('issue-uuid-1'); + expect((sink.submitted[0].event.payload as any).identifier).toBe('ENG-1234'); + expect((sink.submitted[1].event.payload as any).identifier).toBe('ENG-1235'); + }); + + it('refuses when LINEAR_API_TOKEN is not set (fail-closed on missing auth)', async () => { + const sink = recordingSink(); + // Clear env for this test only. Vitest per-file workers make this safe. + const prior = process.env['LINEAR_API_TOKEN']; + delete process.env['LINEAR_API_TOKEN']; + try { + await expect(pollLinearOnce({}, sink, { + fetcher: async () => RECORDED_ISSUES, + })).rejects.toThrow(/LINEAR_API_TOKEN/); + } finally { + if (prior !== undefined) process.env['LINEAR_API_TOKEN'] = prior; + } + expect(sink.submitted).toHaveLength(0); + }); + + it('surfaces GraphQL errors (fail-closed on server rejection)', async () => { + const sink = recordingSink(); + await expect(pollLinearOnce({}, sink, { + token: 't', + fetcher: async () => JSON.stringify({ + errors: [{ message: 'Not authenticated', extensions: { code: 'AUTHENTICATION_ERROR' } }], + }), + })).rejects.toThrow(/Linear GraphQL error/); + expect(sink.submitted).toHaveLength(0); + }); + + it('refuses malformed responses rather than submitting nothing silently', async () => { + const sink = recordingSink(); + await expect(pollLinearOnce({}, sink, { + token: 't', + fetcher: async () => '{"data":{"issues":{"wrongkey":[]}}}', + })).rejects.toThrow(/missing data\.issues\.nodes/); + }); + + it('passes createdAfter filter into the GraphQL query when provided', async () => { + const sink = recordingSink(); + let capturedQuery: string | undefined; + await pollLinearOnce({}, sink, { + token: 't', + createdAfter: '2026-08-31T00:00:00Z', + fetcher: async (query) => { + capturedQuery = query; + return JSON.stringify({ data: { issues: { nodes: [] } } }); + }, + }); + expect(capturedQuery).toBeDefined(); + expect(capturedQuery).toContain('filter:'); + expect(capturedQuery).toContain('createdAt: { gt: "2026-08-31T00:00:00Z" }'); + }); + + it('exports the GraphQL endpoint constant for humans', () => { + expect(LINEAR_GRAPHQL_ENDPOINT).toBe('https://api.linear.app/graphql'); + }); +}); diff --git a/testdata/linear-monitor.flow.yaml b/testdata/linear-monitor.flow.yaml new file mode 100644 index 000000000..6baa2fdc5 --- /dev/null +++ b/testdata/linear-monitor.flow.yaml @@ -0,0 +1,41 @@ +version: '0.1.0' +name: linear-monitor +description: >- + Analyze newly-created Linear issues for triage priority. Second proactive + workload on gate 2 primitives — proves the pattern generalizes beyond + hn-monitor. +triggers: + - id: linear-issue-created + executor: agent-worker + eventType: linear.issue_created + pattern: + type: issue + dedupeKeyTemplate: '{{event.type}}:{{payload.id}}' +steps: + - id: triage-issue + type: agent + instruction: >- + Analyze the triggering Linear issue from the wake context. Determine: + priority (P0/P1/P2/P3), estimated size (S/M/L/XL), and one-line + recommended first action. Output a JSON summary. + recoveryMode: reset + verification: + type: json_schema + schema: + type: object + required: + - issue_id + - priority + - size + - first_action + properties: + issue_id: + type: string + priority: + type: string + enum: [P0, P1, P2, P3] + size: + type: string + enum: [S, M, L, XL] + first_action: + type: string