Skip to content
Open
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
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ All notable changes to this project will be documented in this file.
## [Unreleased]

### Fixed
- Posting an agent activity to Linear now costs one API request instead of three. Both activity-posting paths read the newly created activity back through `AgentActivityPayload.agentActivity`, an unmemoized Linear SDK getter that issues a fresh query on every property access, and both touched it twice while consuming only its `id`. The two extra requests per activity tripled Cyrus's consumption of the Linear API budget and drove rate limiting on busy sessions; a read-back that failed also left an unhandled promise rejection behind. Activity ids now come from `agentActivityId`, which the mutation response already carries, falling back to `agentActivity` for the CLI issue-tracker adapter, whose payload holds the id there as an already-resolved promise and carries no `agentActivityId` at all. The fallback costs nothing on the Linear path, where `agentActivityId` is always present and short-circuits before the getter is reached. ([#1448](https://github.com/cyrusagents/cyrus/pull/1448))
- EdgeWorker state saves are now atomic, preventing a process interrupted during a save from leaving a truncated state file that strands in-flight sessions; empty and legacy-truncated state files also recover cleanly. Thanks @connor-tembo for the contribution. ([CYPACK-1486](https://linear.app/ceedar/issue/CYPACK-1486/can-you-add-a-changelog-entry-for-this), [#1444](https://github.com/cyrusagents/cyrus/pull/1444))

### Changed
Expand Down
23 changes: 19 additions & 4 deletions packages/edge-worker/src/ActivityPoster.ts
Original file line number Diff line number Diff line change
Expand Up @@ -29,10 +29,25 @@ export class ActivityPoster {
try {
const result = await issueTracker.createAgentActivity(input);
if (result.success) {
if (result.agentActivity) {
const activity = await result.agentActivity;
this.logger.debug(`Created ${label} activity ${activity.id}`);
return activity.id;
// `result.agentActivity` is an unmemoized getter on the Linear SDK
// payload: every access constructs a fresh AgentActivityQuery and
// issues another round trip to read back the activity we just
// created. Only the id was ever consumed, and `result.agentActivityId`
// returns it from the mutation response the SDK is already holding, at
// no request cost.
//
// The fallback is for CLIIssueTrackerService, whose payload carries the
// id only on `agentActivity`, as an already-resolved promise property,
// and omits `agentActivityId` entirely. `??` short-circuits on the SDK
// path, so the getter is never touched there; on the CLI adapter it
// awaits a settled promise, which is not a request. Keeping the getter
// out of the `if` is what stops an unawaited read-back from being
// orphaned when it rejects.
const activityId =
result.agentActivityId ?? (await result.agentActivity)?.id;
if (activityId) {
this.logger.debug(`Created ${label} activity ${activityId}`);
return activityId;
}
this.logger.debug(`Created ${label}`);
return null;
Expand Down
23 changes: 20 additions & 3 deletions packages/edge-worker/src/sinks/LinearActivitySink.ts
Original file line number Diff line number Diff line change
Expand Up @@ -97,9 +97,26 @@ export class LinearActivitySink implements IActivitySink {
}),
});

if (result.success && result.agentActivity) {
const agentActivity = await result.agentActivity;
return { activityId: agentActivity.id };
// `result.agentActivity` is an unmemoized getter on the Linear SDK
// payload: every access constructs a fresh AgentActivityQuery and issues
// another round trip to read back the activity we just created. Only the
// id was ever consumed, and `result.agentActivityId` returns it from the
// mutation response the SDK is already holding, at no request cost.
//
// The fallback is for CLIIssueTrackerService, whose payload carries the id
// only on `agentActivity`, as an already-resolved promise property, and
// omits `agentActivityId` entirely. `??` short-circuits on the SDK path, so
// the getter is never touched there; on the CLI adapter it awaits a settled
// promise, which is not a request. Keeping the getter out of the `if` is
// what stops an unawaited read-back from being orphaned when it rejects,
// and reading it only once `success` holds keeps failed mutations free.
if (result.success) {
const activityId =
result.agentActivityId ?? (await result.agentActivity)?.id;

if (activityId) {
return { activityId };
}
}

return {};
Expand Down
159 changes: 159 additions & 0 deletions packages/edge-worker/test/ActivityPoster.request-count.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,159 @@
/**
* Regression tests for the number of Linear requests that one
* ActivityPoster.postActivityDirect() call issues.
*
* postActivityDirect is the second path that posts an agent activity, and it
* read the created activity back through the same unmemoized SDK getter as
* LinearActivitySink did. Its try/catch does not make the orphaned read-back
* safe: the promise produced by the `if` condition is never awaited, so its
* rejection is not the catch block's to handle.
*/

import type { IIssueTrackerService, ILogger } from "cyrus-core";
import { beforeEach, describe, expect, it, vi } from "vitest";
import { ActivityPoster } from "../src/ActivityPoster.js";
import {
countingCreateAgentActivity,
createAgentActivityPayload,
createCLIAdapterAgentActivityPayload,
createRequestLog,
type RequestLog,
} from "./agent-activity-payload-double.js";

describe("ActivityPoster request count", () => {
let poster: ActivityPoster;
let issueTracker: IIssueTrackerService;
let createAgentActivity: ReturnType<typeof vi.fn>;
let logger: ILogger;
let log: RequestLog;

const input = {
agentSessionId: "session-1",
content: { type: "thought" as const, body: "Analyzing..." },
};

const respondWith = (payload: object) => {
createAgentActivity.mockImplementation(
countingCreateAgentActivity(log, payload),
);
};

beforeEach(() => {
log = createRequestLog();
createAgentActivity = vi.fn();
issueTracker = { createAgentActivity } as unknown as IIssueTrackerService;
logger = {
debug: vi.fn(),
error: vi.fn(),
warn: vi.fn(),
info: vi.fn(),
} as unknown as ILogger;

poster = new ActivityPoster(
new Map([["workspace-1", issueTracker]]),
new Map(),
logger,
);
});

it("should issue exactly one Linear request per posted activity", async () => {
respondWith(createAgentActivityPayload(log, "activity-1"));

const activityId = await poster.postActivityDirect(
issueTracker,
input,
"thought",
);

expect(activityId).toBe("activity-1");
expect(log.operations).toEqual(["mutation:agentActivityCreate"]);
});

it("should not read the activity back after creating it", async () => {
respondWith(createAgentActivityPayload(log, "activity-2"));

await poster.postActivityDirect(issueTracker, input, "thought");

expect(log.operations).not.toContain("query:agentActivity");
});

it("should not orphan a rejected read-back promise", async () => {
respondWith(
createAgentActivityPayload(log, "activity-3", { readBackRejects: true }),
);

const unhandled: unknown[] = [];
const onUnhandledRejection = (reason: unknown) => unhandled.push(reason);
process.on("unhandledRejection", onUnhandledRejection);

try {
await poster.postActivityDirect(issueTracker, input, "thought");
// Yield a macrotask so Node can report any rejection left unhandled.
await new Promise((resolve) => setTimeout(resolve, 10));
} finally {
process.off("unhandledRejection", onUnhandledRejection);
}

expect(unhandled).toEqual([]);
});

it("should still return null and log an error when the mutation did not succeed", async () => {
respondWith({ success: false, lastSyncId: 1 });

const activityId = await poster.postActivityDirect(
issueTracker,
input,
"thought",
);

expect(activityId).toBeNull();
expect(logger.error).toHaveBeenCalled();
expect(log.operations).toEqual(["mutation:agentActivityCreate"]);
});

it("should not touch the read-back getter when the id the payload carries is empty", async () => {
// Distinguishes `??` from `||`. Both short-circuit on a populated id, so
// only a falsy-but-present id tells them apart: `??` keeps it and skips
// the getter, `||` falls through and pays for a read-back.
respondWith(createAgentActivityPayload(log, ""));

const activityId = await poster.postActivityDirect(
issueTracker,
input,
"thought",
);

expect(activityId).toBeNull();
expect(log.operations).toEqual(["mutation:agentActivityCreate"]);
});

it("should return the id when the payload comes from the CLI issue-tracker adapter", async () => {
// CLIIssueTrackerService carries the id only on `agentActivity`, as an
// already-resolved promise property, and omits `agentActivityId`. Reading
// `agentActivityId` alone drops the id on that adapter.
respondWith(createCLIAdapterAgentActivityPayload("activity-cli"));

const activityId = await poster.postActivityDirect(
issueTracker,
input,
"thought",
);

expect(activityId).toBe("activity-cli");
// Awaiting an already-resolved promise is not a round trip.
expect(log.operations).toEqual(["mutation:agentActivityCreate"]);
});

it("should still return null when the payload carries no activity id", async () => {
respondWith({ success: true, lastSyncId: 1, agentActivityId: undefined });

const activityId = await poster.postActivityDirect(
issueTracker,
input,
"thought",
);

expect(activityId).toBeNull();
expect(logger.error).not.toHaveBeenCalled();
});
});
148 changes: 148 additions & 0 deletions packages/edge-worker/test/LinearActivitySink.request-count.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,148 @@
/**
* Regression tests for the number of Linear requests that one
* LinearActivitySink.postActivity() call issues.
*
* These live in their own file because they need a different test double from
* the one in LinearActivitySink.test.ts: the payload has to reproduce the SDK's
* unmemoized `agentActivity` getter before the cost of touching it twice is
* observable at all. See ./agent-activity-payload-double.ts.
*/

import type { AgentActivityContent, IIssueTrackerService } from "cyrus-core";
import { beforeEach, describe, expect, it, vi } from "vitest";
import { LinearActivitySink } from "../src/sinks/LinearActivitySink.js";
import {
countingCreateAgentActivity,
createAgentActivityPayload,
createCLIAdapterAgentActivityPayload,
createRequestLog,
type RequestLog,
} from "./agent-activity-payload-double.js";

describe("LinearActivitySink request count", () => {
let sink: LinearActivitySink;
let mockIssueTracker: IIssueTrackerService;
let log: RequestLog;

const mockWorkspaceId = "workspace-123";
const mockSessionId = "session-456";

const activity: AgentActivityContent = {
type: "thought",
body: "Analyzing the codebase...",
};

const respondWith = (payload: object) => {
vi.mocked(mockIssueTracker.createAgentActivity).mockImplementation(
countingCreateAgentActivity(log, payload),
);
};

beforeEach(() => {
log = createRequestLog();
mockIssueTracker = {
createAgentActivity: vi.fn(),
createAgentSessionOnIssue: vi.fn(),
} as unknown as IIssueTrackerService;

sink = new LinearActivitySink(mockIssueTracker, mockWorkspaceId);
});

it("should issue exactly one Linear request per posted activity", async () => {
respondWith(createAgentActivityPayload(log, "activity-1"));

const result = await sink.postActivity(mockSessionId, activity);

expect(result).toEqual({ activityId: "activity-1" });
expect(log.operations).toEqual(["mutation:agentActivityCreate"]);
expect(log.operations).toHaveLength(1);
});

it("should not read the activity back after creating it", async () => {
respondWith(createAgentActivityPayload(log, "activity-2"));

await sink.postActivity(mockSessionId, activity);

expect(log.operations).not.toContain("query:agentActivity");
});

it("should issue one request per activity across a burst of posts", async () => {
respondWith(createAgentActivityPayload(log, "activity-3"));

await sink.postActivity(mockSessionId, activity);
await sink.postActivity(mockSessionId, activity);
await sink.postActivity(mockSessionId, activity);

expect(log.operations).toHaveLength(3);
});

it("should not orphan a rejected read-back promise", async () => {
// The getter access in the `if` condition produced a promise that nobody
// awaited. When that read-back failed -- which is exactly what happens
// once the extra requests have exhausted the rate-limit budget -- its
// rejection was unhandled.
respondWith(
createAgentActivityPayload(log, "activity-4", { readBackRejects: true }),
);

const unhandled: unknown[] = [];
const onUnhandledRejection = (reason: unknown) => unhandled.push(reason);
process.on("unhandledRejection", onUnhandledRejection);

try {
await sink.postActivity(mockSessionId, activity).catch(() => undefined);
// Yield a macrotask so Node can report any rejection left unhandled.
await new Promise((resolve) => setTimeout(resolve, 10));
} finally {
process.off("unhandledRejection", onUnhandledRejection);
}

expect(unhandled).toEqual([]);
});

it("should still return an empty result when the mutation did not succeed", async () => {
respondWith({ success: false, lastSyncId: 1 });

const result = await sink.postActivity(mockSessionId, activity);

expect(result).toEqual({});
expect(log.operations).toEqual(["mutation:agentActivityCreate"]);
});

it("should not touch the read-back getter when the id the payload carries is empty", async () => {
// Distinguishes `??` from `||`. Both short-circuit on a populated id, so
// only a falsy-but-present id tells them apart: `??` keeps it and skips
// the getter, `||` falls through and pays for a read-back.
respondWith(createAgentActivityPayload(log, ""));

const result = await sink.postActivity(mockSessionId, activity);

expect(result).toEqual({});
expect(log.operations).toEqual(["mutation:agentActivityCreate"]);
});

it("should return the id when the payload comes from the CLI issue-tracker adapter", async () => {
// CLIIssueTrackerService carries the id only on `agentActivity`, as an
// already-resolved promise property, and omits `agentActivityId`. Reading
// `agentActivityId` alone drops the id on that adapter.
respondWith(createCLIAdapterAgentActivityPayload("activity-cli"));

const result = await sink.postActivity(mockSessionId, activity);

expect(result).toEqual({ activityId: "activity-cli" });
// Awaiting an already-resolved promise is not a round trip.
expect(log.operations).toEqual(["mutation:agentActivityCreate"]);
});

it("should still return an empty result when the payload carries no activity id", async () => {
respondWith({
success: true,
lastSyncId: 1,
agentActivityId: undefined,
});

const result = await sink.postActivity(mockSessionId, activity);

expect(result).toEqual({});
});
});
Loading