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
Original file line number Diff line number Diff line change
Expand Up @@ -142,7 +142,6 @@ export class HostDailyReviewCoordinator {
this.#closeTask ??= (async () => {
this.beginDrain();
await Promise.allSettled([...this.#inFlight.values()].map((entry) => entry.promise));
this.#store.close();
})();
return this.#closeTask;
}
Expand Down
189 changes: 39 additions & 150 deletions packages/runtime-host/src/server/execution-composition.ts
Original file line number Diff line number Diff line change
Expand Up @@ -48,37 +48,16 @@ import {
} from '@maka/runtime/shell-detect';
import { type MakaTool } from '@maka/runtime/tool-runtime';
import { type RuntimeHostedRootAuthority } from '@maka/runtime/message-authority';
import {
openInteractiveProjectCatalogForWrite,
type InteractiveProjectCatalogWriter,
} from '@maka/storage';
import { createAgentGraphControlStore } from '@maka/storage/agent-graph-control-store';
import {
createArtifactAttachmentResourceReader,
createReadImageSnapshotter,
openInteractiveArtifactStoreForWrite,
} from '@maka/storage/artifact-stores';
import { openInteractiveScheduledTaskStoreForWrite } from '@maka/storage/scheduled-task-store';
import { openInteractiveDeepResearchStoreForWrite } from '@maka/storage/deep-research-authority';
import { openInteractiveDailyReviewAuthorityForWrite } from '@maka/storage/daily-review-authority';
import { openInteractiveGoalAuthorityForWrite } from '@maka/storage/goal-authority';
import { openInteractivePlanStoreForWrite } from '@maka/storage/plan-authority';
import {
isSessionNotFoundError,
openInteractiveExecutionStoresForWrite,
} from '@maka/storage/execution-stores';
import { isSessionNotFoundError } from '@maka/storage/execution-stores';
import { createExternalSessionAdapterRegistry } from '@maka/storage/external-sessions';
import { createGitWorktreeChildExecutor } from '@maka/storage/git-worktree-child-executor';
import {
type InteractiveLongTermMemoryWriter,
openInteractiveLongTermMemoryStoreForWrite,
} from '@maka/storage/long-term-memory-store';
import { openInteractiveMemoryBundleStoreForWrite } from '@maka/storage/memory-bundle-store';
import { runWithStorageRootLease } from '@maka/storage/root-authority';
import { openInteractiveRuntimePolicyStoresForWrite } from '@maka/storage/runtime-policy-stores';
import { openInteractiveTaskLedgerStoreForWrite } from '@maka/storage/task-ledger-authority';
import { openInteractiveShellRunStoreForWrite } from '@maka/storage/shell-run-authority';
import { openInteractiveUsageStoresForWrite } from '@maka/storage/usage-stores';
import { openStorageWriterComposition } from '@maka/storage';
import { resolveWorkspaceIdentity } from '@maka/storage/workspace-identity';
import {
openManagedWorkspaceOwner,
Expand Down Expand Up @@ -214,78 +193,45 @@ export async function createExecutionRuntimeHostComposition(
options: CreateExecutionRuntimeHostCompositionOptions = {},
dependencies: ExecutionRuntimeHostCompositionDependencies = {},
): Promise<ExecutionRuntimeHostComposition> {
const stores = await openInteractiveExecutionStoresForWrite(context.owner.lease);
await stores.sessionStore.ready();
const storage = await openStorageWriterComposition(context.owner.lease, {
afterRuntimePolicyOpened: async (stores) => {
if (options.bootstrapRuntimePolicy !== false) {
await ensureBootstrapRuntimePolicy({
workspaceRoot: context.owner.capability.canonicalPath,
stores,
onDeferredError: (error) =>
console.error(
`[runtime-host] optional bootstrap target could not be configured: ${generalizedErrorMessage(error)}`,
),
});
}
},
});
const stores = storage.execution;
let graphControlStore: ReturnType<typeof createAgentGraphControlStore> | undefined;
let taskLedgerStore:
| Awaited<ReturnType<typeof openInteractiveTaskLedgerStoreForWrite>>
| undefined;
let usageStores: Awaited<ReturnType<typeof openInteractiveUsageStoresForWrite>> | undefined;
let artifactStore: Awaited<ReturnType<typeof openInteractiveArtifactStoreForWrite>> | undefined;
let shellRunStore: Awaited<ReturnType<typeof openInteractiveShellRunStoreForWrite>> | undefined;
let longTermMemoryStore: InteractiveLongTermMemoryWriter | undefined;
let scheduledTaskStore:
| Awaited<ReturnType<typeof openInteractiveScheduledTaskStoreForWrite>>
| undefined;
let planStore: Awaited<ReturnType<typeof openInteractivePlanStoreForWrite>> | undefined;
let deepResearchStore:
| Awaited<ReturnType<typeof openInteractiveDeepResearchStoreForWrite>>
| undefined;
let dailyReviewStore:
| Awaited<ReturnType<typeof openInteractiveDailyReviewAuthorityForWrite>>
| undefined;
let goalStore: Awaited<ReturnType<typeof openInteractiveGoalAuthorityForWrite>> | undefined;
let graphClient: HostAgentGraphCoordinator | undefined;
let sessionEffects: HostSessionEffectCoordinator | undefined;
let memoryExtraction: HostMemoryExtractionCoordinator | undefined;
let unsubscribeTaskLedger: (() => void) | undefined;
let unsubscribeTranscriptChanges: (() => void) | undefined;
let managedWorkspaceOwner: ManagedWorkspaceOwner | undefined;
let workspaceExecution: RuntimeHostWorkspaceExecutionComposition | undefined;
let projectCatalog: InteractiveProjectCatalogWriter | undefined;
let goalExecutions: HostGoalExecutionCoordinator | undefined;
try {
const openedProjectCatalog = await openInteractiveProjectCatalogForWrite(context.owner.lease);
projectCatalog = openedProjectCatalog;
const runtimePolicyStores = await openInteractiveRuntimePolicyStoresForWrite(
context.owner.lease,
);
if (options.bootstrapRuntimePolicy !== false) {
await ensureBootstrapRuntimePolicy({
workspaceRoot: context.owner.capability.canonicalPath,
stores: runtimePolicyStores,
onDeferredError: (error) =>
console.error(
`[runtime-host] optional bootstrap target could not be configured: ${generalizedErrorMessage(error)}`,
),
});
}
const openedProjectCatalog = storage.projectCatalog;
const runtimePolicyStores = storage.runtimePolicy;
const oauthCredentials = new HostOAuthExecutionAuthority(runtimePolicyStores);
const openedScheduledTaskStore = await openInteractiveScheduledTaskStoreForWrite(
context.owner.lease,
);
scheduledTaskStore = openedScheduledTaskStore;
const openedPlanStore = await openInteractivePlanStoreForWrite(context.owner.lease);
planStore = openedPlanStore;
const openedDeepResearchStore = await openInteractiveDeepResearchStoreForWrite(
context.owner.lease,
);
deepResearchStore = openedDeepResearchStore;
const openedDailyReviewStore = await openInteractiveDailyReviewAuthorityForWrite(
context.owner.lease,
);
dailyReviewStore = openedDailyReviewStore;
const openedGoalStore = await openInteractiveGoalAuthorityForWrite(context.owner.lease);
goalStore = openedGoalStore;
const memoryStore = await openInteractiveMemoryBundleStoreForWrite(context.owner.lease);
longTermMemoryStore = await openInteractiveLongTermMemoryStoreForWrite(context.owner.lease);
taskLedgerStore = await openInteractiveTaskLedgerStoreForWrite(context.owner.lease);
const openedArtifactStore = await openInteractiveArtifactStoreForWrite(context.owner.lease);
artifactStore = openedArtifactStore;
const openedUsageStores = await openInteractiveUsageStoresForWrite(context.owner.lease);
usageStores = openedUsageStores;
const openedShellRunStore = await openInteractiveShellRunStoreForWrite(context.owner.lease);
shellRunStore = openedShellRunStore;
const openedScheduledTaskStore = storage.scheduledTasks;
const openedPlanStore = storage.plan;
const openedDeepResearchStore = storage.deepResearch;
const openedDailyReviewStore = storage.dailyReview;
const openedGoalStore = storage.goal;
const memoryStore = storage.memoryBundle;
const longTermMemoryStore = storage.longTermMemory;
const taskLedgerStore = storage.taskLedger;
const openedArtifactStore = storage.artifacts;
const openedUsageStores = storage.usage;
const openedShellRunStore = storage.shellRuns;
const worktreeChildExecutor = createGitWorktreeChildExecutor({
storageRoot: context.owner.capability.canonicalPath,
});
Expand Down Expand Up @@ -1398,24 +1344,18 @@ export async function createExecutionRuntimeHostComposition(
state: () => requireMemory(memory).recover(),
},
drain: [() => memoryExtraction?.beginDrain(), () => memory?.beginDrain()],
close: [
() => memoryExtraction?.close(),
() => memory?.close(),
() => longTermMemoryStore?.close(),
],
close: [() => memoryExtraction?.close(), () => memory?.close()],
releaseConnection: [
(connectionId) => requireMemory(memory).releaseConnection(connectionId),
],
}),
createRuntimeHostDomainModule({
id: 'plan',
handlers: [plans.handlers],
close: [() => openedPlanStore.close()],
}),
createRuntimeHostDomainModule({
id: 'project-catalog',
handlers: [projects.handlers],
close: [() => openedProjectCatalog.close()],
}),
createRuntimeHostDomainModule({
id: 'session',
Expand All @@ -1437,7 +1377,6 @@ export async function createExecutionRuntimeHostComposition(
}
},
},
close: [() => stores.sessionStore.close?.()],
}),
createRuntimeHostDomainModule({
id: 'configuration',
Expand Down Expand Up @@ -1469,14 +1408,10 @@ export async function createExecutionRuntimeHostComposition(
() => (backendInvalidationPoisoned ? undefined : manager.refreshIdleBackends()),
() => skills.close(),
() => oauth?.close(),
() => openedUsageStores.close(),
() => openedArtifactStore.close(),
() => {
unsubscribeTranscriptChanges?.();
unsubscribeTaskLedger?.();
taskLedgerStore?.close();
},
() => shellRunStore?.close(),
],
releaseConnection: [(connectionId) => artifacts.releaseConnection(connectionId)],
}),
Expand All @@ -1497,7 +1432,7 @@ export async function createExecutionRuntimeHostComposition(
createRuntimeHostDomainModule({
id: 'deep-research',
handlers: [requireDeepResearch(deepResearch).handlers],
close: [() => deepResearch?.close(), () => openedDeepResearchStore.close()],
close: [() => deepResearch?.close()],
}),
createRuntimeHostDomainModule({
id: 'daily-review',
Expand Down Expand Up @@ -1596,7 +1531,7 @@ export async function createExecutionRuntimeHostComposition(
domains: () => requireGoal(goal).recover(),
},
drain: [() => goalExecutions?.beginDrain(), () => goal?.beginDrain()],
close: [() => requireGoal(goal).close(), () => openedGoalStore.close()],
close: [() => requireGoal(goal).close()],
}),
createRuntimeHostDomainModule({
id: 'session-retirement',
Expand Down Expand Up @@ -1633,6 +1568,11 @@ export async function createExecutionRuntimeHostComposition(
} catch (error) {
errors.push(error);
}
try {
await storage.close();
} catch (error) {
errors.push(error);
}
if (poisonFailure && !errors.includes(poisonFailure)) errors.push(poisonFailure);
if (errors.length > 0) {
throw new AggregateError(errors, 'Unable to close Runtime Host execution composition');
Expand Down Expand Up @@ -1668,16 +1608,6 @@ export async function createExecutionRuntimeHostComposition(
} catch (closeError) {
errors.push(closeError);
}
try {
dailyReviewStore?.close();
} catch (closeError) {
errors.push(closeError);
}
try {
deepResearchStore?.close();
} catch (closeError) {
errors.push(closeError);
}
try {
graphClient?.close();
} catch (closeError) {
Expand All @@ -1688,25 +1618,9 @@ export async function createExecutionRuntimeHostComposition(
} catch (closeError) {
errors.push(closeError);
}
try {
await usageStores?.close();
} catch (closeError) {
errors.push(closeError);
}
try {
artifactStore?.close();
} catch (closeError) {
errors.push(closeError);
}
try {
unsubscribeTranscriptChanges?.();
unsubscribeTaskLedger?.();
taskLedgerStore?.close();
} catch (closeError) {
errors.push(closeError);
}
try {
shellRunStore?.close();
} catch (closeError) {
errors.push(closeError);
}
Expand All @@ -1716,32 +1630,7 @@ export async function createExecutionRuntimeHostComposition(
errors.push(closeError);
}
try {
longTermMemoryStore?.close();
} catch (closeError) {
errors.push(closeError);
}
try {
scheduledTaskStore?.close();
} catch (closeError) {
errors.push(closeError);
}
try {
await goalStore?.close();
} catch (closeError) {
errors.push(closeError);
}
try {
planStore?.close();
} catch (closeError) {
errors.push(closeError);
}
try {
projectCatalog?.close();
} catch (closeError) {
errors.push(closeError);
}
try {
await stores.sessionStore.close?.();
await storage.close();
} catch (closeError) {
errors.push(closeError);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -212,7 +212,6 @@ export class HostScheduledTaskCoordinator implements ScheduledTaskToolAuthority
this.beginDrain();
await this.#lane;
this.#closed = true;
this.#store.close();
}

async create(input: {
Expand Down
1 change: 1 addition & 0 deletions packages/storage/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@
"./pet-pack-store": "./dist/pet-pack-store.js",
"./root-authority": "./dist/root-authority.js",
"./state-root-composition": "./dist/state-root-composition.js",
"./storage-writer-composition": "./dist/storage-writer-composition.js",
"./runtime-policy-stores": "./dist/runtime-policy-stores.js",
"./scheduled-task-store": "./dist/scheduled-task-store.js",
"./settings-store": "./dist/settings-store.js",
Expand Down
Loading