Skip to content

Commit 98ed911

Browse files
author
linyuan.yang
committed
ACP Service
1 parent 8f3a53a commit 98ed911

10 files changed

Lines changed: 478 additions & 386 deletions

File tree

packages/channel.xiaoai/tsconfig.tsbuildinfo

Lines changed: 1 addition & 1 deletion
Large diffs are not rendered by default.
0 Bytes
Binary file not shown.

packages/sbot/src/Agent/ACPAgentPool.ts

Lines changed: 27 additions & 32 deletions
Original file line numberDiff line numberDiff line change
@@ -1,10 +1,10 @@
1-
import { ACPAgentService } from "scorpio.ai";
1+
import { PersistentACPAgentService } from "scorpio.ai";
22
import { LoggerService } from "../Core/LoggerService";
33

44
const logger = LoggerService.getLogger('ACPAgentPool');
55

66
export interface ACPPoolEntry {
7-
instance: ACPAgentService;
7+
instance: PersistentACPAgentService;
88
key: string;
99
agentId: string;
1010
agentName: string;
@@ -35,53 +35,48 @@ export class ACPAgentPool {
3535
return ACPAgentPool.inst;
3636
}
3737

38-
async acquire(
39-
key: string,
40-
agentId: string,
41-
agentName: string,
42-
dbSessionId: string,
43-
configHash: string,
44-
factory: () => Promise<ACPAgentService>,
45-
): Promise<ACPAgentService> {
38+
async tryGet(key: string, configHash: string): Promise<PersistentACPAgentService | null> {
4639
const existing = this.pool.get(key);
40+
if (!existing) return null;
4741

48-
if (existing) {
49-
if (existing.configHash !== configHash) {
50-
logger.info(`Config changed for ${key}, recreating`);
51-
await existing.instance.forceDispose();
52-
this.pool.delete(key);
53-
} else if (existing.instance.isAlive()) {
54-
existing.lastAccessed = Date.now();
55-
return existing.instance;
56-
} else {
57-
logger.warn(`Dead process for ${key}, recreating`);
58-
this.pool.delete(key);
59-
}
42+
if (existing.configHash !== configHash) {
43+
logger.info(`Config changed for ${key}, recreating`);
44+
await existing.instance.forceDispose();
45+
this.pool.delete(key);
46+
return null;
47+
}
48+
49+
if (!existing.instance.isAlive()) {
50+
logger.warn(`Dead process for ${key}, recreating`);
51+
this.pool.delete(key);
52+
return null;
6053
}
6154

62-
const instance = await factory();
55+
existing.lastAccessed = Date.now();
56+
return existing.instance;
57+
}
58+
59+
put(key: string, instance: PersistentACPAgentService, meta: { agentId: string; agentName: string; dbSessionId: string; configHash: string }): void {
6360
instance.pooled = true;
64-
instance.onExit = () => {
61+
instance.onPoolExit = () => {
6562
const entry = this.pool.get(key);
6663
if (entry?.instance === instance) {
6764
logger.warn(`ACP process died unexpectedly: ${key}`);
6865
this.pool.delete(key);
6966
}
7067
};
7168

72-
const entry: ACPPoolEntry = {
69+
this.pool.set(key, {
7370
instance,
7471
key,
75-
agentId,
76-
agentName,
77-
dbSessionId,
78-
configHash,
72+
agentId: meta.agentId,
73+
agentName: meta.agentName,
74+
dbSessionId: meta.dbSessionId,
75+
configHash: meta.configHash,
7976
createdAt: Date.now(),
8077
lastAccessed: Date.now(),
81-
};
82-
this.pool.set(key, entry);
78+
});
8379
logger.info(`Cached ACP instance: ${key}`);
84-
return instance;
8580
}
8681

8782
async release(key: string): Promise<void> {

packages/sbot/src/Agent/AgentFactory.ts

Lines changed: 21 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
import {
22
AgentServiceBase, SingleAgentService, GenerativeAgentService,
3-
ACPAgentService, T_ACPCommand, T_ACPArgs, T_ACPEnv, T_ACPSessionMode, T_ACPWorkPath,
3+
TransientACPAgentService, PersistentACPAgentService, T_ACPCommand, T_ACPArgs, T_ACPEnv, T_ACPWorkPath,
44
IModelService,
55
IAgentSaverService, AgentMemorySaver, ILoggerService,
66
ConversationCompactor, IConversationCompactor, T_SummaryModelService, T_CompactPromptTemplate,
@@ -274,37 +274,34 @@ export class AgentFactory {
274274
): Promise<AgentServiceBase> {
275275
const workPath = options.workPath ?? process.cwd();
276276
const sessionMode = entry.sessionMode ?? ACPSessionMode.Persistent;
277+
const acpArgs = {
278+
[T_ACPCommand]: entry.command,
279+
[T_ACPArgs]: entry.args ?? [],
280+
[T_ACPWorkPath]: workPath,
281+
...(entry.env && { [T_ACPEnv]: entry.env }),
282+
};
277283

278284
if (sessionMode === ACPSessionMode.Persistent) {
279285
const pool = ACPAgentPool.getInstance();
280286
const key = `${options.agentId}:${options.dbSessionId}`;
281287
const agentName = entry.name ?? options.agentId;
282288
const configHash = JSON.stringify({ command: entry.command, args: entry.args ?? [], env: entry.env ?? {}, workPath });
283289

284-
return pool.acquire(key, options.agentId, agentName, options.dbSessionId, configHash, async () => {
285-
const sub = new ServiceContainer();
286-
if (container.isRegistered(ILoggerService))
287-
sub.registerInstance(ILoggerService, await container.resolve(ILoggerService));
288-
if (container.isRegistered(IAgentSaverService))
289-
sub.registerInstance(IAgentSaverService, await container.resolve(IAgentSaverService));
290-
sub.registerWithArgs(ACPAgentService, {
291-
[T_ACPCommand]: entry.command,
292-
[T_ACPArgs]: entry.args ?? [],
293-
[T_ACPWorkPath]: workPath,
294-
...(entry.env && { [T_ACPEnv]: entry.env }),
295-
[T_ACPSessionMode]: ACPSessionMode.Persistent,
296-
});
297-
return sub.resolve<ACPAgentService>(ACPAgentService);
298-
});
290+
const cached = await pool.tryGet(key, configHash);
291+
if (cached) return cached;
292+
293+
const sub = new ServiceContainer();
294+
if (container.isRegistered(ILoggerService))
295+
sub.registerInstance(ILoggerService, await container.resolve(ILoggerService));
296+
if (container.isRegistered(IAgentSaverService))
297+
sub.registerInstance(IAgentSaverService, await container.resolve(IAgentSaverService));
298+
sub.registerWithArgs(PersistentACPAgentService, acpArgs);
299+
const instance = await sub.resolve<PersistentACPAgentService>(PersistentACPAgentService);
300+
pool.put(key, instance, { agentId: options.agentId, agentName, dbSessionId: options.dbSessionId, configHash });
301+
return instance;
299302
}
300303

301-
container.registerWithArgs(ACPAgentService, {
302-
[T_ACPCommand]: entry.command,
303-
[T_ACPArgs]: entry.args ?? [],
304-
[T_ACPWorkPath]: workPath,
305-
...(entry.env && { [T_ACPEnv]: entry.env }),
306-
[T_ACPSessionMode]: ACPSessionMode.Transient,
307-
});
308-
return container.resolve(ACPAgentService);
304+
container.registerWithArgs(TransientACPAgentService, acpArgs);
305+
return container.resolve(TransientACPAgentService);
309306
}
310307
}

0 commit comments

Comments
 (0)