Skip to content
Merged
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
19 changes: 19 additions & 0 deletions apps/server/src/orchestration/Layers/CheckpointReactor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@ import {
} from "../../checkpointing/Utils.ts";
import * as CheckpointStore from "../../checkpointing/CheckpointStore.ts";
import { ProviderService } from "../../provider/Services/ProviderService.ts";
import { PrimeAgentRecoveryLedger } from "../../provider/prime/PrimeAgentRecoveryLedger.ts";
import { CheckpointReactor, type CheckpointReactorShape } from "../Services/CheckpointReactor.ts";
import { forkParked } from "../../serverActivation.ts";
import { OrchestrationEngineService } from "../Services/OrchestrationEngine.ts";
Expand Down Expand Up @@ -85,6 +86,9 @@ const make = Effect.gen(function* () {
const orchestrationEngine = yield* OrchestrationEngineService;
const projectionSnapshotQuery = yield* ProjectionSnapshotQuery;
const providerService = yield* ProviderService;
const recoveryLedger = Option.getOrUndefined(
yield* Effect.serviceOption(PrimeAgentRecoveryLedger),
);
const checkpointStore = yield* CheckpointStore.CheckpointStore;
const receiptBus = yield* RuntimeReceiptBus;
const workspaceEntries = yield* WorkspaceEntries.WorkspaceEntries;
Expand Down Expand Up @@ -823,6 +827,21 @@ const make = Effect.gen(function* () {
),
),
);
if (recoveryLedger !== undefined) {
yield* Effect.gen(function* () {
yield* recoveryLedger.markCheckpointQuiesced({
threadId: event.threadId,
updatedAt: yield* nowIso,
});
yield* recoveryLedger.deleteIfSettled(event.threadId);
}).pipe(
Effect.catchCause(() =>
Effect.logWarning("failed to settle Prime Agent recovery checkpoint proof", {
threadId: event.threadId,
}),
),
);
}
return;
}
});
Expand Down
22 changes: 22 additions & 0 deletions apps/server/src/orchestration/Layers/ProviderRuntimeIngestion.ts
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,7 @@ import { makeDrainableWorker } from "@t3tools/shared/DrainableWorker";

import { ProviderRegistry } from "../../provider/Services/ProviderRegistry.ts";
import { ProviderService } from "../../provider/Services/ProviderService.ts";
import { PrimeAgentRecoveryLedger } from "../../provider/prime/PrimeAgentRecoveryLedger.ts";
import {
rateLimitFromRuntimeEventPayload,
usageWindowsFromRuntimeEventPayload,
Expand Down Expand Up @@ -1507,8 +1508,23 @@ const make = Effect.gen(function* () {
const orchestrationEngine = yield* OrchestrationEngineService;
const projectionSnapshotQuery = yield* ProjectionSnapshotQuery;
const providerService = yield* ProviderService;
const recoveryLedger = Option.getOrUndefined(
yield* Effect.serviceOption(PrimeAgentRecoveryLedger),
);
const providerRegistry = yield* ProviderRegistry;
const projectionTurnRepository = yield* ProjectionTurnRepository;
const settleRecoveryTerminalProjection = (threadId: ThreadId, updatedAt: string) =>
recoveryLedger === undefined
? Effect.void
: recoveryLedger.markTerminalProjected({ threadId, updatedAt }).pipe(
Effect.andThen(recoveryLedger.deleteIfSettled(threadId)),
Effect.catchCause(() =>
Effect.logWarning("failed to settle Prime Agent terminal projection proof", {
threadId,
}),
),
Effect.asVoid,
);
const serverSettingsService = yield* ServerSettingsService;
const providerCommandId = (event: ProviderRuntimeEvent, tag: string) =>
crypto.randomUUIDv4.pipe(
Expand Down Expand Up @@ -2228,6 +2244,9 @@ const make = Effect.gen(function* () {
if ((cleared.eventCount ?? 0) > 0 && stoppedLineageStillProjected) {
yield* clearTurnStateForSession(thread.id);
}
if ((cleared.eventCount ?? 0) > 0) {
yield* settleRecoveryTerminalProjection(thread.id, event.createdAt);
}
return;
}
const eventMatchesStoppedSession =
Expand Down Expand Up @@ -2566,6 +2585,9 @@ const make = Effect.gen(function* () {
"thread-session-set",
);
if ((applied.eventCount ?? 0) === 0) return;
if (event.type === "turn.completed" || event.type === "session.exited") {
yield* settleRecoveryTerminalProjection(thread.id, event.createdAt);
}
}

if (event.type === "turn.started" && acceptedTurnStartedSourcePlan !== null) {
Expand Down
2 changes: 2 additions & 0 deletions apps/server/src/persistence/Migrations.ts
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,7 @@ import Migration0046 from "./Migrations/046_ProjectionThreadLinkedPullRequest.ts
import Migration0047 from "./Migrations/047_ProjectionThreadsUnsettledAt.ts";
import Migration0048 from "./Migrations/048_ProjectionThreadSessionPendingTurnRequest.ts";
import Migration0049 from "./Migrations/049_ProjectionThreadSessionPendingStop.ts";
import Migration0050 from "./Migrations/050_PrimeAgentRecoveryLedger.ts";
/**
* Migration loader with all migrations defined inline.
*
Expand Down Expand Up @@ -144,6 +145,7 @@ export const migrationEntries = [
[47, "ProjectionThreadsUnsettledAt", Migration0047],
[48, "ProjectionThreadSessionPendingTurnRequest", Migration0048],
[49, "ProjectionThreadSessionPendingStop", Migration0049],
[50, "PrimeAgentRecoveryLedger", Migration0050],
] as const;

export const migrationManifest = migrationEntries.map(([id, name]) => [id, name] as const);
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,56 @@
import * as Effect from "effect/Effect";
import * as SqlClient from "effect/unstable/sql/SqlClient";

export default Effect.gen(function* () {
const sql = yield* SqlClient.SqlClient;

yield* sql`
CREATE TABLE IF NOT EXISTS prime_agent_recovery_ledger (
thread_id TEXT PRIMARY KEY,
provider_instance_id TEXT NOT NULL,
session_incarnation_id TEXT NOT NULL,
admission_request_id TEXT NOT NULL,
turn_id TEXT,
package_root TEXT NOT NULL,
package_version TEXT NOT NULL,
managed_build_id TEXT NOT NULL,
sdk_features_json TEXT NOT NULL,
daemon_capabilities_json TEXT NOT NULL,
protocol_name TEXT NOT NULL,
protocol_version INTEGER NOT NULL,
schema_revision INTEGER NOT NULL,
active_session_id TEXT NOT NULL,
native_session_id TEXT NOT NULL,
recovery_handle TEXT NOT NULL,
supervisor_generation TEXT NOT NULL,
ownership_generation INTEGER NOT NULL,
cursor_generation TEXT NOT NULL,
cursor_sequence INTEGER NOT NULL,
correlation_id TEXT NOT NULL,
mcp_owner_id TEXT NOT NULL,
recovery_config_json TEXT NOT NULL,
launch_environment_json TEXT NOT NULL,
transcript_message_count INTEGER NOT NULL DEFAULT 0,
transcript_fingerprints_json TEXT NOT NULL DEFAULT '[]',
owner_token TEXT NOT NULL,
state TEXT NOT NULL,
adoption_previous_owner_token TEXT,
adoption_owner_token TEXT,
adoption_request_id TEXT,
adoption_mcp_owner_id TEXT,
adoption_phase TEXT,
adoption_attempt INTEGER NOT NULL DEFAULT 0,
adoption_recovery_handle TEXT,
adoption_proof_json TEXT,
native_cleanup_proven INTEGER NOT NULL DEFAULT 0,
terminal_projected INTEGER NOT NULL DEFAULT 0,
checkpoint_quiesced INTEGER NOT NULL DEFAULT 0,
updated_at TEXT NOT NULL
)
`;

yield* sql`
CREATE INDEX IF NOT EXISTS idx_prime_agent_recovery_ledger_active
ON prime_agent_recovery_ledger(state, provider_instance_id, updated_at)
`;
});
38 changes: 37 additions & 1 deletion apps/server/src/provider/Drivers/PrimeAgentDriver.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,12 +4,13 @@ import {
type ServerProvider,
type ServerProviderDistribution,
} from "@t3tools/contracts";
import { HostProcessPlatform } from "@t3tools/shared/hostProcess";
import { HostProcessArchitecture, HostProcessPlatform } from "@t3tools/shared/hostProcess";
import { resolveCommandPath } from "@t3tools/shared/shell";
import * as Crypto from "effect/Crypto";
import * as Effect from "effect/Effect";
import * as FileSystem from "effect/FileSystem";
import * as Path from "effect/Path";
import * as Result from "effect/Result";
import * as Schema from "effect/Schema";
import { HttpClient } from "effect/unstable/http";
import { ChildProcessSpawner } from "effect/unstable/process";
Expand Down Expand Up @@ -123,6 +124,7 @@ export const PrimeAgentDriver: ProviderDriver<PrimeAgentSettings, PrimeAgentDriv
create: ({ instanceId, displayName, accentColor, environment, enabled, config }) =>
Effect.gen(function* () {
const hostPlatform = yield* HostProcessPlatform;
const hostArchitecture = yield* HostProcessArchitecture;
if (!isPrimeAgentProviderPlatformSupported(hostPlatform)) {
return yield* new ProviderDriverError({
driver: DRIVER_KIND,
Expand Down Expand Up @@ -219,6 +221,37 @@ export const PrimeAgentDriver: ProviderDriver<PrimeAgentSettings, PrimeAgentDriv
env: processEnv,
});

const recoveryDistribution = yield* Effect.result(
Effect.gen(function* () {
const executablePath = path.resolve(
yield* resolveCommandPath(effectiveConfig.binaryPath || "prime-agent", {
env: processEnv,
}),
);
const publicPackage = yield* locatePrimeAgentPublicPackage(executablePath);
const distribution = yield* Effect.promise(() =>
inspectPrimeAgentDistribution(
{
stateDir: serverConfig.stateDir,
instanceId,
packageRoot: publicPackage.packageRoot,
platform: hostPlatform,
checkedAt: "1970-01-01T00:00:00.000Z",
enableUpdateChecks: false,
},
{ loadLatestVerifiedPublication },
),
);
return { publicPackage, distribution };
}),
);
const recoveryManagedBuildId =
Result.isSuccess(recoveryDistribution) &&
recoveryDistribution.success.distribution.classification === "pylon-managed" &&
recoveryDistribution.success.distribution.buildId !== null
? recoveryDistribution.success.distribution.buildId
: undefined;

const backend = yield* negotiatePrimeAgentBackend(
{
enabled: effectiveConfig.enabled,
Expand All @@ -228,6 +261,8 @@ export const PrimeAgentDriver: ProviderDriver<PrimeAgentSettings, PrimeAgentDriv
environment: processEnv,
stateDir: serverConfig.stateDir,
providerInstanceId: instanceId,
recoveryEnabled: recoveryManagedBuildId !== undefined,
architecture: hostArchitecture,
},
{
resolveExecutable: (command, resolvedEnvironment) =>
Expand Down Expand Up @@ -383,6 +418,7 @@ export const PrimeAgentDriver: ProviderDriver<PrimeAgentSettings, PrimeAgentDriv
environment: processEnv,
...(eventLoggers.native ? { nativeEventLogger: eventLoggers.native } : {}),
instanceId,
...(recoveryManagedBuildId === undefined ? {} : { recoveryManagedBuildId }),
} as const;
const adapter =
backend.runtime === "daemon"
Expand Down
Loading
Loading