Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
46 commits
Select commit Hold shift + click to select a range
34037fa
refactor(deploy): declare credentials that stay on the edge
matanyall Sep 17, 2026
b8a427e
feat(linear): persist app installations and incoming agent events
matanyall Sep 17, 2026
6cf0793
fix(linear): declare local runtime files and plan references
matanyall Sep 17, 2026
c47cb7e
feat(linear): add native session conversations and authenticated edge…
matanyall Sep 17, 2026
52aef8f
feat(linear): dispatch native sessions and recover durable deliveries
matanyall Sep 17, 2026
25e8379
feat(linear): acknowledge sessions durably before dispatch
matanyall Sep 17, 2026
9cbae92
feat(linear): share private files through native session uploads
matanyall Sep 17, 2026
556711c
feat(linear): manage issues within the requesting person’s access
matanyall Sep 17, 2026
5b17864
feat(linear): let requesters stop their own session work
matanyall Sep 17, 2026
c774739
test(review): budget time for the real Git diff fixture
matanyall Sep 17, 2026
ce287fa
fix(linear): keep sessions working after follow-up acknowledgement
matanyall Sep 17, 2026
cb817e5
feat(linear): run the bot without Slack credentials
matanyall Sep 17, 2026
204abdc
fix(linear): report missing channel configuration without a stack
matanyall Sep 17, 2026
1130aff
fix(linear): run shared-session follow-ups with their own requester
matanyall Sep 17, 2026
dcce0aa
refactor(linear): share current requester access facts
matanyall Sep 17, 2026
a219fdd
fix(linear): recheck requester access before dispatch and recovery
matanyall Sep 17, 2026
800599a
feat(linear): ask for input without completing the session
matanyall Sep 17, 2026
459e6a5
fix(linear): authorize cancellation of waiting questions
matanyall Sep 17, 2026
d6be6ce
refactor(channels): share inbound attachment classification
matanyall Sep 17, 2026
7eb9461
feat(linear): read private files in prompts and history
matanyall Sep 17, 2026
a043766
chore(linear): merge main and preserve clarification replies
matanyall Sep 17, 2026
840652e
fix(linear): classify inbound files in the edge Worker
matanyall Sep 17, 2026
6a3132d
fix(linear): preserve mention files and finish edge downloads
matanyall Sep 17, 2026
ba5d37a
fix(linear): wait for answers from final agent turns
matanyall Sep 17, 2026
f4197a2
fix(linear): separate local intake and run-page origins
matanyall Sep 17, 2026
dd02c69
feat(linear): open native sessions for delegated child work
matanyall Sep 18, 2026
7ef6052
chore(linear): merge main and reconcile decision numbering
matanyall Sep 18, 2026
4c281f5
fix(linear): preserve requester isolation in child replies
matanyall Sep 18, 2026
fdeca77
fix(linear): keep child questions pending in run tools
matanyall Sep 18, 2026
18193ba
fix(linear): hold coordinator rounds while children need input
matanyall Sep 18, 2026
851f391
test(linear): type coordinator question fixtures
matanyall Sep 18, 2026
784b0ba
chore(linear): merge main while preserving clarification guards
matanyall Sep 18, 2026
d2d7dc8
fix(linear): resume coordinated work after clarification replies
matanyall Sep 18, 2026
dac7209
fix(linear): retry unavailable coordinator clarification state
matanyall Sep 18, 2026
8f86851
chore: merge main into Linear integration
matanyall Sep 18, 2026
7ff9e94
test(linear): cover refusal rendering after main integration
matanyall Sep 18, 2026
4815deb
fix(linear): persist cancellation of waiting questions
matanyall Sep 18, 2026
01d9062
fix(linear): bound coordinator waits and retry history failures
matanyall Sep 18, 2026
63e2222
refactor(linear): share authorized file context lookup
matanyall Sep 18, 2026
7564de6
feat(linear): stage private files in agent workspaces
matanyall Sep 18, 2026
8d5bcc9
fix(linear): cancel incoming file transfers on hard stop
matanyall Sep 18, 2026
d349fc0
test(execution): wait for the restarted runtime log
matanyall Sep 18, 2026
b137f25
fix(linear): cancel queued requests before dispatch
matanyall Sep 18, 2026
991e793
fix(linear): preserve stop through deferred admission
matanyall Sep 18, 2026
2bce684
fix(linear): await delivery admission before execution
matanyall Sep 18, 2026
b0a3c0b
feat(linear): support project and document session origins
matanyall Sep 18, 2026
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
7 changes: 7 additions & 0 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,13 @@ E2B_API_KEY=e2b_...
# PORT=8080
# PUBLIC_BASE_URL=https://switchboard.example.com

# Linear-only or combined startup (see docs/how-to/connect-linear.md).
# Only the shared bridge token belongs in the bot process; OAuth secrets stay at the edge.
# LINEAR_BRIDGE_TOKEN=<random-shared-bearer>
# Optional separate edge origin, e.g. local OAuth/intake on :8080 and bot run pages on :8082.
# Defaults to PUBLIC_BASE_URL. Use HTTPS, or HTTP only on loopback.
# LINEAR_BRIDGE_URL=http://localhost:8080

# GitHub identity for the coding/review agents — pick ONE:
# (a) GitHub App (idiomatic: org-owned, no user, 1h tokens minted on demand).
# Org settings -> Developer settings -> GitHub Apps -> New GitHub App;
Expand Down
3 changes: 3 additions & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,9 @@ data/
workspaces/
.env
.agent-env/
.wrangler/
.dev.vars
.dev.vars.*
config/config.yaml
node_modules
build.json
Expand Down
2 changes: 2 additions & 0 deletions config/config.example.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -202,6 +202,8 @@ defaults:
# slack:U0456DEV:
# actions: [agent:run:coding, repo:write, friction:write]
# repos: [acme/api] # the repos this user may use
# linear:*:
# actions: [work-items:write] # issue edits still check each person's live Linear access
# access:svc:ops-bot: # an Access service token (its common_name)
# actions: [runs:read, runs:write, friction:read]
# channels: all # sees runs from every channel
Expand Down
25 changes: 25 additions & 0 deletions deploy/cloudflare-memory/runs.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -218,6 +218,31 @@ describe("run usage", () => {
});

describe("run history routes", () => {
it("persists a waiting Stop and fences late history writes and stale stops", async () => {
const key = storeKey();
const now = Date.now();
const question = record("question", now, { awaitingInput: true });
const stop = { at: now + 1, mode: "hard", by: { kind: "chat", id: "linear:org:alice" } };
await post("/runs/put", { storeKey: key, record: question });
expect(
(await post("/runs/stop-waiting", { storeKey: key, id: "question", stop: { ...stop, at: now - 1 } })).data,
).toEqual({ result: "conflict" });
expect((await post("/runs/stop-waiting", { storeKey: key, id: "question", stop })).data).toEqual({
result: "stopped",
});
await post("/runs/put", { storeKey: key, record: question });
const read = await post("/runs/get", { storeKey: key, id: "question" });
expect(read.data.record).toMatchObject({ status: "stopped_hard", inputStop: stop });
expect((read.data.record as RunRecord).awaitingInput).toBeUndefined();
expect((await post("/runs/stop-waiting", { storeKey: key, id: "question", stop })).data).toEqual({
result: "stopped",
});
expect(
(await post("/runs/stop-waiting", { storeKey: key, id: "question", stop: { ...stop, by: {} } })).status,
).toBe(400);
expect((await post("/runs/stop-waiting", { storeKey: key, id: "question", stop }, {})).status).toBe(401);
});

it("put → get round-trips the record with events in seq order; unknown id → {record: null} 200", async () => {
const key = storeKey();
const now = Date.now();
Expand Down
53 changes: 47 additions & 6 deletions deploy/cloudflare-memory/worker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,10 @@ import {
type DeliverySnapshotPatch,
} from "../../src/core/deliverySnapshotStore.ts";
import {
type InputStop,
isInputStop,
preserveInputStop,
stopWaitingRecord,
applyRetention,
clampRetentionPolicy,
isRunRecord,
Expand Down Expand Up @@ -2308,6 +2312,14 @@ export class RunHistoryDO extends DurableObject<Env> {
{
const now = systemClock();
const policy = proposal ? this.applyProposal(proposal, now).policy : this.policyState().policy;
const existing = this.sql
.exec<RunRow>(
`SELECT run_id, agent, channel_id, finished_at, bytes, event_count, summary_json FROM runs WHERE run_id = ?`,
record.id,
)
.toArray()[0];
const previous = existing ? parseSummary(existing) : undefined;
record = preserveInputStop(record, previous?.inputStop);
const finishedAt = Math.min(record.finishedAt, now + RUN_MAX_FUTURE_MS);
// The tracing stamps get the same skew clamp (docs/reference/specs/tracing.md).
const stored: RunRecord = {
Expand All @@ -2320,12 +2332,6 @@ export class RunHistoryDO extends DurableObject<Env> {
};
const { events, ...summary } = stored;
const bytes = utf8ByteLength(JSON.stringify(stored));
const existing = this.sql
.exec<{ event_count: number; finished_at: number; bytes: number }>(
`SELECT event_count, finished_at, bytes FROM runs WHERE run_id = ?`,
record.id,
)
.toArray()[0];
const unchanged =
existing !== undefined &&
sameStoredVersion(
Expand Down Expand Up @@ -2409,6 +2415,33 @@ export class RunHistoryDO extends DurableObject<Env> {
}
}

/** Persist cancellation before a channel closes the waiting conversation. */
async stopWaiting(id: string, stop: InputStop): Promise<"stopped" | "conflict" | "not_found"> {
let result: "stopped" | "conflict" | "not_found" = "not_found";
let changed: RunRecord | undefined;
this.ctx.storage.transactionSync(() => {
const row = this.sql
.exec<RunRow>(
`SELECT run_id, agent, channel_id, finished_at, bytes, event_count, summary_json FROM runs WHERE run_id = ?`,
id,
)
.toArray()[0];
if (!row || !this.isKept(row, this.policyState().policy, systemClock())) return;
const summary = parseSummary(row);
if (!summary) return;
const record = { ...summary, events: parseEventRows(this.eventRows(id, 0, Number.MAX_SAFE_INTEGER)) };
changed = stopWaitingRecord(record, stop);
if (!changed) {
result = "conflict";
return;
}
this.upsertInTransaction(changed);
result = "stopped";
});
if (changed) await sendRunFinished(this.env.SHIP_COORDINATOR, changed);
return result;
}

/** Remove a run and its events. Returns whether a run row existed. */
async delete(id: string): Promise<boolean> {
let deleted = false;
Expand Down Expand Up @@ -4241,6 +4274,13 @@ async function handleRuns(pathname: string, body: unknown, env: Env): Promise<Re
);
return json(result);
}
if (pathname === "/runs/stop-waiting") {
const parsed = parseRunTarget(body);
if (!parsed.ok) return json({ error: parsed.error }, 400);
const stop = (body as Record<string, unknown>).stop;
if (!isInputStop(stop)) return json({ error: "invalid input stop" }, 400);
return json({ result: await stub(parsed.value.storeKey).stopWaiting(parsed.value.id, stop) });
}
if (pathname === "/runs/get") {
const parsed = parseRunTarget(body);
if (!parsed.ok) return json({ error: parsed.error }, 400);
Expand Down Expand Up @@ -4301,6 +4341,7 @@ const ROUTES = new Set([
"/schedules/record",
"/schedules/latest",
"/runs/put",
"/runs/stop-waiting",
"/runs/get",
"/runs/summary",
"/runs/list",
Expand Down
13 changes: 13 additions & 0 deletions deploy/cloudflare/linear.local.jsonc
Original file line number Diff line number Diff line change
@@ -0,0 +1,13 @@
{
"$schema": "node_modules/wrangler/config-schema.json",
"name": "switchboard-linear-local",
"main": "linear.local.ts",
"compatibility_date": "2026-08-01",
"compatibility_flags": ["enable_request_signal"],
"workers_dev": false,
"vars": { "PUBLIC_BASE_URL": "http://localhost:8080" },
"durable_objects": {
"bindings": [{ "name": "LINEAR_STATE", "class_name": "LinearState" }]
},
"migrations": [{ "tag": "v1", "new_sqlite_classes": ["LinearState"] }]
}
10 changes: 10 additions & 0 deletions deploy/cloudflare/linear.local.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,10 @@
// Local OAuth and webhook testing needs no container, Docker image or cloud account.
import { handleLinearEdge, linearRoute, type LinearEnv } from "./linear";
export { LinearState } from "./linear";

export default {
fetch(request: Request, env: LinearEnv): Promise<Response> | Response {
if (!linearRoute(new URL(request.url).pathname)) return new Response("Not found", { status: 404 });
return handleLinearEdge(request, env);
},
};
209 changes: 209 additions & 0 deletions deploy/cloudflare/linear.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,209 @@
import { DurableObject } from "cloudflare:workers";
import {
LinearOAuth,
LinearTokenProvider,
LINEAR_AUTHORIZE_PATH,
LINEAR_CALLBACK_PATH,
} from "../../src/channels/linear/oauth.js";
import { DirectLinearApi } from "../../src/channels/linear/api.js";
import { handleLinearBridge, LINEAR_BRIDGE_PATH } from "../../src/channels/linear/bridge.js";
import { StoredLinearStore, type LinearOAuthState } from "../../src/channels/linear/store.js";
import { StoredLinearChildStore } from "../../src/channels/linear/children.js";
import { SqlLinearInbox } from "../../src/channels/linear/inbox.js";
import { LinearAcknowledgements } from "../../src/channels/linear/acknowledgement.js";
import { revokeLinearInstallation } from "../../src/channels/linear/lifecycle.js";
import { boundedBody, handleLinearWebhook, LINEAR_WEBHOOK_PATH } from "../../src/channels/linear/webhook.js";
import { LINEAR_TIMING } from "../../src/core/budgets.js";
import { systemClock } from "../../src/core/trace/clock.js";
import type { Env } from "./worker";

export type LinearEnv = Pick<
Env,
| "LINEAR_STATE"
| "LINEAR_CLIENT_ID"
| "LINEAR_CLIENT_SECRET"
| "LINEAR_APPLICATION_ID"
| "LINEAR_WEBHOOK_SECRET"
| "LINEAR_ORGANIZATION_ID"
| "LINEAR_BRIDGE_TOKEN"
| "PUBLIC_BASE_URL"
| "ARTIFACTS"
>;

/** Public OAuth and signed webhook routes never wake the bot's container.
* An installation without the optional credentials stays explicitly disabled. */
export function linearRoute(pathname: string): boolean {
return (
pathname === LINEAR_AUTHORIZE_PATH ||
pathname === LINEAR_CALLBACK_PATH ||
pathname === LINEAR_WEBHOOK_PATH ||
pathname === LINEAR_BRIDGE_PATH
);
}

export async function handleLinearEdge(request: Request, env: LinearEnv): Promise<Response> {
if (
!env.LINEAR_CLIENT_ID ||
!env.LINEAR_CLIENT_SECRET ||
!env.LINEAR_APPLICATION_ID ||
!env.LINEAR_WEBHOOK_SECRET ||
!env.PUBLIC_BASE_URL
) {
return Response.json({ error: "linear_disabled" }, { status: 503, headers: { "cache-control": "no-store" } });
}
// Buffer only bounded raw bytes before crossing the Durable Object boundary.
// An early rejection there must not leave a streaming subrequest pumping
// from the outer request after its response has already been sent.
if (request.body) {
const bytes = await boundedBody(request);
if (!bytes) return Response.json({ error: "too_large" }, { status: 413 });
request = new Request(request, { body: bytes as Uint8Array<ArrayBuffer> });
}
return env.LINEAR_STATE.get(env.LINEAR_STATE.idFromName("installation")).fetch(request, {
signal: request.signal,
});
}

/** One durable host for the installation's OAuth state, credentials and inbox.
* No credential route exists: tokens never leave this object's storage through
* an HTTP response. The queue is retained independently of container lifetimes. */
export class LinearState extends DurableObject<LinearEnv> {
private readonly store: StoredLinearStore;
private readonly inbox: SqlLinearInbox;
private readonly tokens: LinearTokenProvider;
private readonly acknowledgements: LinearAcknowledgements;

constructor(ctx: DurableObjectState, env: LinearEnv) {
super(ctx, env);
this.store = new StoredLinearStore(ctx.storage);
this.inbox = new SqlLinearInbox(ctx.storage.sql);
this.tokens = new LinearTokenProvider({
clientId: env.LINEAR_CLIENT_ID ?? "",
clientSecret: env.LINEAR_CLIENT_SECRET ?? "",
store: this.store,
fetch: (input, init) => fetch(input, init),
clock: systemClock,
});
this.acknowledgements = new LinearAcknowledgements({
inbox: this.inbox,
api: (organizationId) => this.api(organizationId),
clock: systemClock,
warn: (message) => console.warn(message),
});
}

private async api(organizationId: string): Promise<DirectLinearApi> {
if (this.env.LINEAR_ORGANIZATION_ID && organizationId !== this.env.LINEAR_ORGANIZATION_ID)
throw new Error("linear_wrong_installation");
const installation = await this.store.getInstallation(organizationId);
if (!installation) throw new Error("linear_not_installed");
return new DirectLinearApi({
organizationId,
appUserId: installation.appUserId,
children: new StoredLinearChildStore(this.ctx.storage),
...(this.env.ARTIFACTS
? {
copy: {
put: (key: string, stream: ReadableStream<Uint8Array>, type: string) =>
this.env.ARTIFACTS!.put(key, stream, { httpMetadata: { contentType: type } }),
lengthPipe: (size: number) => new FixedLengthStream(size),
},
}
: {}),
token: () => this.tokens.accessToken(organizationId),
fetch: (input, init) => fetch(input, init),
});
}

private async armAlarm(delay: number): Promise<void> {
const next = systemClock() + delay;
const current = await this.ctx.storage.getAlarm();
if (current === null || current > next) await this.ctx.storage.setAlarm(next);
}

async fetch(request: Request): Promise<Response> {
const env = this.env;
if (
!env.LINEAR_CLIENT_ID ||
!env.LINEAR_CLIENT_SECRET ||
!env.LINEAR_APPLICATION_ID ||
!env.LINEAR_WEBHOOK_SECRET ||
!env.PUBLIC_BASE_URL
) {
return Response.json({ error: "linear_disabled" }, { status: 503 });
}
try {
if ((await this.ctx.storage.getAlarm()) === null)
await this.ctx.storage.setAlarm(systemClock() + LINEAR_TIMING.oauthStateMs);
const path = new URL(request.url).pathname;
if (path === LINEAR_BRIDGE_PATH)
return await handleLinearBridge(request, {
token: env.LINEAR_BRIDGE_TOKEN,
inbox: this.inbox,
clock: systemClock,
api: (organizationId) => this.api(organizationId),
});
if (path === LINEAR_WEBHOOK_PATH)
return await handleLinearWebhook(request, {
secret: env.LINEAR_WEBHOOK_SECRET,
applicationId: env.LINEAR_APPLICATION_ID,
organizationId: env.LINEAR_ORGANIZATION_ID,
clock: systemClock,
accept: async (event) => {
if (event.payload.type === "OAuthApp" && event.payload.action === "revoked") {
if (!(await revokeLinearInstallation(this.store, event))) return false;
await this.inbox.cancelOrganization(event.payload.organizationId, event.receivedAt);
}
const acknowledge = event.payload.type === "AgentSessionEvent" && event.payload.action === "created";
const accepted = await this.inbox.accept(event, { acknowledge });
if (acknowledge) {
// The alarm is durable before HTTP 200; network work runs outside
// the response lifetime and never waits for the bot container.
await this.armAlarm(LINEAR_TIMING.progressMs);
this.ctx.waitUntil(
this.acknowledgements.flush().catch(() => {
console.warn("[linear] acknowledgement sweep unavailable; alarm will retry");
}),
);
}
return accepted;
},
});
if (path === LINEAR_AUTHORIZE_PATH) {
await this.pruneStates();
const pending = await this.ctx.storage.list({ prefix: "oauth:", limit: 256 });
if (pending.size >= 256) return Response.json({ error: "too_many_pending_installations" }, { status: 429 });
}
const oauth = new LinearOAuth({
clientId: env.LINEAR_CLIENT_ID,
clientSecret: env.LINEAR_CLIENT_SECRET,
baseUrl: env.PUBLIC_BASE_URL,
organizationId: env.LINEAR_ORGANIZATION_ID,
store: this.store,
fetch: (input, init) => fetch(input, init),
clock: systemClock,
});
return await oauth.handle(request);
} catch {
return Response.json({ error: "linear_unavailable" }, { status: 503 });
}
}

private async pruneStates(): Promise<void> {
const now = systemClock();
const states = await this.ctx.storage.list<LinearOAuthState>({ prefix: "oauth:" });
const expired = [...states].filter(([, value]) => value.expiresAt <= now).map(([key]) => key);
for (let offset = 0; offset < expired.length; offset += 128)
await this.ctx.storage.delete(expired.slice(offset, offset + 128));
}

async alarm(): Promise<void> {
try {
await this.acknowledgements.flush();
await this.pruneStates();
await this.inbox.prune(systemClock() - LINEAR_TIMING.deliveryRetentionMs);
} finally {
await this.armAlarm((await this.inbox.hasPendingAcks()) ? LINEAR_TIMING.progressMs : LINEAR_TIMING.oauthStateMs);
}
}
}
Loading