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
51 changes: 45 additions & 6 deletions packages/codingcode/src/agent/agent.ts
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,10 @@ import { isTurnEnd } from '../sink/types.js';
import type { SessionRef } from '../session/types.js';
import type { PermissionMode } from '../util/enums.js';
import type { ToolCatalog, ToolResult } from '../tools/types.js';
import { mediaKindOf, textOf, textPart, type IncomingPart } from '../llm/types.js';
import type { ToolCall } from '../llm/types.js';
import { capabilitiesOf } from '../infra/models.js';
import { sniffMediaMime } from '../util/media.js';
import { loadConfig } from '../infra/config.js';
import { createLogger } from '../infra/logger.js';
import { normalizePath } from '../util/path.js';
Expand All @@ -42,6 +45,35 @@ function toolOutcomeOf(result: ToolResult): ToolOutcome {
: { status: result.status, output: result.output };
}

function assertMediaAllowed(
input: readonly IncomingPart[],
model: string
): Effect.Effect<void, AgentError> {
return Effect.gen(function* () {
let needsVision = false;
let needsAudio = false;
for (const p of input) {
if (p.type !== 'media') continue;
const mimeType = sniffMediaMime(p.bytes) ?? p.declaredMimeType ?? '';
if (mediaKindOf(mimeType) === 'audio') needsAudio = true;
else needsVision = true;
}
if (!needsVision && !needsAudio) return;

// capabilitiesOf 是同步查询,抛出的 AgentError 原样进错误通道
const caps = yield* Effect.try({
try: () => capabilitiesOf(model),
catch: (e) => (e instanceof AgentError ? e : AgentError.invalidInput(String(e))),
});
if (needsVision && !caps.vision) {
return yield* Effect.fail(AgentError.invalidInput('model does not accept image or PDF input'));
}
if (needsAudio && !caps.audio) {
return yield* Effect.fail(AgentError.invalidInput('model does not accept audio input'));
}
});
}

const logger = createLogger();

function toFrameError(e: AgentError): FrameError {
Expand Down Expand Up @@ -81,7 +113,7 @@ export const AgentLayer = Layer.effect(
)
);

const runTurn = (input: string, opts: RunTurnOptions) =>
const runTurn = (input: IncomingPart[], opts: RunTurnOptions) =>
Effect.gen(function* () {
const normalizedCwd = normalizePath(opts.cwd);

Expand All @@ -106,7 +138,7 @@ export const AgentLayer = Layer.effect(
normalizedCwd,
{
model,
title: input,
title: textOf(input) || 'New session',
activeProfile: opts.activeProfile,
permissionMode: opts.permissionMode,
},
Expand Down Expand Up @@ -148,7 +180,9 @@ export const AgentLayer = Layer.effect(

const toolEnv = yield* toolEnvPort.getToolEnv();

const turnId = (yield* session.recordUser(state, input)).turnId;
yield* assertMediaAllowed(input, model);
const parts = yield* session.materializeInput(state, input);
const turnId = (yield* session.recordUser(state, parts)).turnId;

// 用户显式 @ 的 skill:按 path 回查权威数据,正文拼块后作为同回合的第二条 user 事件
if (opts.skills?.length) {
Expand All @@ -158,7 +192,7 @@ export const AgentLayer = Layer.effect(
const entries = yield* Effect.forEach(chosen, (s) =>
skills.readContent(s.skillPath).pipe(Effect.map((body) => ({ skill: s, body })))
);
yield* session.recordSystem(state, renderSkillBlock(entries));
yield* session.recordSystem(state, [textPart(renderSkillBlock(entries))]);
}
}

Expand Down Expand Up @@ -347,7 +381,12 @@ export const AgentLayer = Layer.effect(
Effect.tryPromise({
try: async () => {
for await (const part of llm.completeStream(
{ messages: llmMessages, system, tools, maxSteps: 1 },
{
messages: llmMessages,
system,
tools,
maxSteps: 1,
},
model,
abortSignal
)) {
Expand Down Expand Up @@ -436,7 +475,7 @@ export const AgentLayer = Layer.effect(
}
stopContinuations++;
const injection = stopDecision.injection ?? '(continue)';
const systemEv = yield* session.recordSystem(state, injection);
const systemEv = yield* session.recordSystem(state, [textPart(injection)]);
yield* context.absorb(sessionRef, [systemEv]);
continue;
}
Expand Down
3 changes: 2 additions & 1 deletion packages/codingcode/src/agent/port.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ import { Context } from 'effect';
import type { Effect } from 'effect';
import type { FrameBody } from '../sink/types.js';
import type { ProfileName, PermissionMode } from '../util/enums.js';
import type { IncomingPart } from '../llm/types.js';

import type { AgentError } from '../util/error.js';

Expand All @@ -20,7 +21,7 @@ export interface RunTurnOptions {

export interface AgentShape {
runTurn(
input: string,
input: IncomingPart[],
opts: RunTurnOptions
): Effect.Effect<
{
Expand Down
64 changes: 52 additions & 12 deletions packages/codingcode/src/context/context.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,8 @@ import { Layer, Effect } from 'effect';
import { randomUUID } from 'crypto';
import { loadConfig } from '../infra/config.js';
import { transcriptPathOf } from '../session/paths.js';
import type { Message } from '../llm/types.js';
import type { Message, ResolvedMessage, ResolvedContentPart, ResolvedMediaPart } from '../llm/types.js';
import { textOf, textPart } from '../llm/types.js';
import type {
SessionEvent,
AssistantEvent,
Expand Down Expand Up @@ -111,7 +112,7 @@ export function buildContextMessages(
break;
case 'assistant': {
const ev = event as AssistantEvent;
const msg: Message = { role: 'assistant', content: event.content };
const msg: Message = { role: 'assistant', content: [textPart(event.content)] };
if (event.toolCalls && event.toolCalls.length > 0) {
msg.tool_calls = event.toolCalls.map((tc) => ({
id: tc.id,
Expand All @@ -135,17 +136,21 @@ export function buildContextMessages(
resolvedIds.add(event.toolCallId);
messages.push({
role: 'tool',
content: output,
content: [textPart(output)],
tool_call_id: event.toolCallId,
tool_name: event.toolName,
});
break;
}
case 'summary':
messages.push({ role: 'system', name: 'compacted_history', content: event.summaryText });
messages.push({
role: 'system',
name: 'compacted_history',
content: [textPart(event.summaryText)],
});
break;
case 'subagent_result':
messages.push({ role: 'user', content: event.content });
messages.push({ role: 'user', content: [textPart(event.content)] });
break;
}
}
Expand Down Expand Up @@ -180,7 +185,8 @@ export function buildContextMessages(
if (curr.role === prev.role && curr.role !== 'system') {
if (curr.role === 'tool') continue;
if (curr.role === 'assistant' && curr.tool_calls && curr.tool_calls.length > 0) continue;
prev.content += '\n\n' + curr.content;
// parts 数组本身就是分隔符,providers 侧按 part 逐个送出
prev.content = [...prev.content, ...curr.content];
filtered.splice(i, 1);
}
}
Expand Down Expand Up @@ -322,7 +328,11 @@ export const ContextLayer = Layer.effect(
);
buf.events.push(summaryEvent);

const summaryMsg: Message = { role: 'system', name: 'compacted_history', content: summary };
const summaryMsg: Message = {
role: 'system',
name: 'compacted_history',
content: [textPart(summary)],
};
return Math.max(0, totalTokens - estimateMessageTokens(summaryMsg));
});

Expand Down Expand Up @@ -361,15 +371,20 @@ export const ContextLayer = Layer.effect(
const transcriptText = transcript
.map(
(m) =>
`[${m.role}${(m as any).tool_name ? ':' + (m as any).tool_name : ''}]\n${m.content}`
`[${m.role}${(m as any).tool_name ? ':' + (m as any).tool_name : ''}]\n${textOf(m.content)}`
)
.join('\n\n');

const system = COMPACTION_SYSTEM_PROMPT;

const userMsg: Message = {
// 压缩走纯文本投影
const userMsg: ResolvedMessage = {
role: 'user',
content: `Compact the following conversation transcript into the sections above:\n\n${transcriptText}`,
content: [
textPart(
`Compact the following conversation transcript into the sections above:\n\n${transcriptText}`
),
],
};

const result = yield* llm
Expand All @@ -383,7 +398,32 @@ export const ContextLayer = Layer.effect(
return raw.trim();
}

const getHistory = (ref: SessionRef, model: string): Effect.Effect<Message[], AgentError> =>
const resolveMediaParts = (
msgs: Message[],
cwd: string
): Effect.Effect<ResolvedMessage[], AgentError> =>
Effect.gen(function* () {
const assets = new Set<string>();
for (const m of msgs) {
for (const p of m.content) if (p.type === 'media') assets.add(p.asset);
}
if (assets.size === 0) return msgs as ResolvedMessage[];

const resolved = yield* session.resolveAssets(cwd, [...assets]);
return msgs.map((m) => ({
...m,
content: m.content.map((p): ResolvedContentPart => {
if (p.type === 'text') return p;
const dataUrl = resolved.get(p.asset);
if (!dataUrl) return { type: 'text', text: `[media missing: ${p.asset}]` };
const part: ResolvedMediaPart = { type: 'media', dataUrl, mimeType: p.mimeType };
if (p.filename) part.filename = p.filename;
return part;
}),
}));
});

const getHistory = (ref: SessionRef, model: string): Effect.Effect<ResolvedMessage[], AgentError> =>
Effect.gen(function* () {
const buf = yield* ensureBuffer(ref);
const contextWindow = contextWindowOf(model);
Expand All @@ -397,7 +437,7 @@ export const ContextLayer = Layer.effect(
transition: { to: 'executing' },
});
}
return buildContextMessages(buf.events, buf.compactedTurnIds);
return yield* resolveMediaParts(buildContextMessages(buf.events, buf.compactedTurnIds), ref.cwd);
});

const absorb = (ref: SessionRef, events: readonly SessionEvent[]): Effect.Effect<void> =>
Expand Down
6 changes: 3 additions & 3 deletions packages/codingcode/src/context/port.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
import { Context } from 'effect';
import type { Effect } from 'effect';
import type { Message } from '../llm/types.js';
import type { ResolvedMessage } from '../llm/types.js';
import type { SessionRef } from '../session/types.js';
import type { SessionEvent } from '../session/types.js';

Expand All @@ -13,8 +13,8 @@ export interface CompressResult {
}

export interface ContextShape {
/** 取该会话当前给模型的 history。回合内首次调用读一次盘,之后走内存;换回合自动重建 */
getHistory(ref: SessionRef, model: string): Effect.Effect<Message[], AgentError>;
/** 取该会话当前给模型的完整载荷:回合内首次调用读一次盘,之后走内存;换回合自动重建。媒体字节在此装配完成 */
getHistory(ref: SessionRef, model: string): Effect.Effect<ResolvedMessage[], AgentError>;
/** 把本回合自己写进 transcript 的事件并入内存态(零 IO) */
absorb(ref: SessionRef, events: readonly SessionEvent[]): Effect.Effect<void>;
/** 手动压缩入口(HTTP /compact),不受阈值限制 */
Expand Down
39 changes: 37 additions & 2 deletions packages/codingcode/src/context/tokens.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,42 @@
import type { Message } from '../llm/types.js';
import { mediaKindOf, type MediaPart, type Message, type TextPart } from '../llm/types.js';
import type { StoredMediaPart } from '../session/types.js';

/**
* 音频每秒 token 数:取各家公开值里的上界(Gemini 32 / 秒,Qwen-Audio 25 / 秒,
* OpenAI 未公布固定值),让压缩早触发。
*/
const AUDIO_TOKENS_PER_SECOND = 32;
/** PDF 按体积折算:1MB ≈ 5 万 token。 */
const PDF_TOKENS_PER_MB = 50_000;

/**
* 可估算的内容块:瘦块(llm 层)与落盘富块(session 层)都收。
*
* 元数据存在就精算,不存在就取保守回退——估算只做压缩决策,不追求精确。
*/
export type EstimablePart = TextPart | MediaPart | StoredMediaPart;

export function estimateTokensForPart(p: EstimablePart): number {
if (p.type === 'text') return estimateTokensForContent(p.text);
switch (mediaKindOf(p.mimeType)) {
case 'image':
// tile 计费近似:基准 85 + 每 512×512 tile 170
return (
85 +
170 *
Math.ceil((('width' in p ? p.width : undefined) ?? 512) / 512) *
Math.ceil((('height' in p ? p.height : undefined) ?? 512) / 512)
);
case 'audio':
return Math.ceil((('durationSec' in p ? p.durationSec : undefined) ?? 1) * AUDIO_TOKENS_PER_SECOND);
case 'file':
return Math.ceil((('bytes' in p ? p.bytes : 0) / (1024 * 1024)) * PDF_TOKENS_PER_MB);
}
}

export function estimateMessageTokens(m: Message): number {
let tokens = estimateTokensForContent(m.content ?? '');
let tokens = 0;
for (const p of m.content) tokens += estimateTokensForPart(p);
tokens += estimateTokensForContent(m.role);
if (m.name) tokens += estimateTokensForContent(m.name);
if (m.tool_call_id) tokens += estimateTokensForContent(m.tool_call_id);
Expand Down
31 changes: 31 additions & 0 deletions packages/codingcode/src/infra/models.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,18 @@ import { join } from 'path';
import { AgentError } from '../util/error.js';
import { loadConfig, updateActiveModel } from './config.js';

/** config/models.json 里每个模型的 capabilities 声明。 */
export interface ModelCapabilitiesEntry {
vision?: 'supported' | 'unsupported';
audio?: 'supported' | 'unsupported';
}

/** 准入要看的能力位。 */
export interface ModelCapabilities {
vision: boolean;
audio: boolean;
}

/** 模型清单里一条可被会话选中的模型。 */
export interface SelectableModel {
id: string;
Expand All @@ -13,12 +25,14 @@ export interface SelectableModel {
base_url: string;
api_key_env: string;
context_window: number;
capabilities: ModelCapabilities;
}

export interface ModelDescriptor {
id: string;
name: string;
context_window?: number;
capabilities?: ModelCapabilitiesEntry;
}

export interface ProviderEntry {
Expand Down Expand Up @@ -70,6 +84,10 @@ export function flattenModels(cat: ProviderCatalog): SelectableModel[] {
base_url: p.base_url,
api_key_env: p.api_key_env,
context_window: m.context_window ?? DEFAULT_CONTEXT_WINDOW,
capabilities: {
vision: m.capabilities?.vision === 'supported',
audio: m.capabilities?.audio === 'supported',
},
});
}
}
Expand Down Expand Up @@ -115,6 +133,19 @@ export function contextWindowOf(model: string): number {
return entry?.context_window ?? DEFAULT_CONTEXT_WINDOW;
}

/** 指定模型的能力位;模型不存在时报 CONFIG_INVALID。 */
export function capabilitiesOf(model: string): ModelCapabilities {
const target = model?.trim();
if (!target) {
const entry = activeModel();
if (!entry) throw new AgentError('CONFIG_INVALID', activeModelError());
return entry.capabilities;
}
const found = findModel(target);
if (!found) throw new AgentError('CONFIG_INVALID', `Model "${target}" not found in models.json`);
return found.capabilities;
}

export function setGlobalActive(model: string): void {
const found = findModel(model.trim());
if (!found) {
Expand Down
Loading
Loading