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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 10 additions & 0 deletions sdk/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';
149 changes: 149 additions & 0 deletions sdk/src/linear-poller.ts
Original file line number Diff line number Diff line change
@@ -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<unknown>;
}

/** Injected so parsing and submission stay deterministic in tests. */
export type LinearFetcher = (query: string, token: string) => Promise<string>;

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<unknown[]> {
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;
102 changes: 102 additions & 0 deletions sdk/tests/linear-poller.test.ts
Original file line number Diff line number Diff line change
@@ -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');
});
});
41 changes: 41 additions & 0 deletions testdata/linear-monitor.flow.yaml
Original file line number Diff line number Diff line change
@@ -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