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
2 changes: 1 addition & 1 deletion server/modules/providers/list/gjc/GJC-PROVIDER-SPEC.md
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@ Gajae Code App의 provider `gjc`(Gajae Code) 구현 기록. 초기 read-only 세
- `gjc-session-synchronizer.provider.ts` — 원형 `codex-session-synchronizer`. `gjcHome=~/.gjc/agent`, 스캔 `path.join(gjcHome,'sessions')`. `extractFirstValidJsonlData`로 첫 줄 파싱: **`data.id`/`data.cwd` 직접**(codex처럼 payload 아님). title=첫 user message 파생(`extractFirstUserMessageFromStart`를 gjc `type:message,role:user` content-text로 재작성) → 없으면 history.db → 없으면 `Untitled gjc Session`. 파생은 `shared/utils.ts`의 `deriveSessionTitle`(슬래시 커맨드·@멘션·코드 펜스·마크다운 제거, 첫 문장 경계, ≤40자 + `…`)을 거치며, DB에 이미 이름이 있으면 덮어쓰지 않는다. `deriveSessionTitle(filePath)`는 `POST /api/providers/sessions/:id/regenerate-title`이 사용자 요청으로 제목을 다시 파생할 때만 기존 이름을 대체한다. `sessionsDb.createSession(id, 'gjc', cwd, name, createdAt, updatedAt, filePath)`. **모델 제목**: 새 세션의 첫 턴에서 Bun 어댑터가 런타임의 `utils/title-generator`(`generateSessionTitle`, TUI와 동일 조건: 첫 user 메시지·이름 없음·`GJC_NO_TITLE` 미설정)를 호출해 `sessionManager.setSessionName(title,'auto')`로 트랜스크립트 헤더에 기록하고, `{kind:'session_title'}` 메시지를 보낸다. `ChatSessionWriter`가 이를 채팅으로 내보내지 않고 `sessionsDb.applyGeneratedSessionName`으로 저장한 뒤 `session_upserted`로 방송한다. 우선순위는 `sessions.name_source`(`user` | `auto` | `derived` | NULL=구버전 행)로 정한다: 사용자가 지은 이름은 사용자만 바꾸고, 모델 제목은 그 외 전부를 대체하며, 파생 제목은 동기화기가 값을 바꿀 때 찍힌다.
- `server/modules/providers/services/gjc-session-watcher.service.ts` — `gajae-core watch`를 별도 자식 프로세스로 실행해 저장 세션 루트와 live 세션 루트 안에 canonical containment를 통과한 `.jsonl` add/change 이벤트만 64 KiB 제한 NDJSON으로 수신한다. 이벤트는 순서대로 기존 `synchronizeProviderFile('gjc', path)`에 전달하며, 큐 상한·ready 타임아웃·취소 가능한 종료 drain·지수 백오프 재시작·재시작 후 GJC 전용 reconciliation을 적용한다. GJC용 Chokidar fallback은 없고 기존 4개 provider watcher는 그대로 유지한다.
- `gjc-transcript-message.ts` — 일반 메시지와 표시 가능한 사용자 스킬 요청을 공통 해석한다. 히스토리·턴 계보·제목 파생이 공유하며, 스킬의 확장 본문은 반환하지 않는다.
- `gjc-sessions.provider.ts` — `getSessionById(id).jsonl_path` → 제한된 JSONL 스트리밍 → 공통 메시지 해석 → `message.role` + `message.content[]` 파트별 user/assistant/thinking/tool_use/tool_result 정규화. timestamp 정렬 + `sliceTailPage` 페이지네이션(`createNormalizedMessage`/`generateMessageId`, 멀티 text 파트 id 충돌 방지 discriminator).
- `gjc-sessions.provider.ts` — `getSessionById(id).jsonl_path` → 제한된 JSONL 스트리밍 → 공통 메시지 해석 → `message.role` + `message.content[]` 파트별 user/assistant/thinking/tool_use/tool_result 정규화. timestamp 정렬 후 표시 행(tool_result 제외) 기준 tail 페이지네이션: 작은 인덱스 pass 뒤 선택된 행과 그 tool_result만 두 번째 pass에서 payload로 유지하고, 읽는 동안 파일이 바뀌면 재시도 후 409(`HISTORY_CHANGED`), 무제한 요청은 표시 행 5,000개까지만 허용하고 초과 시 413(`HISTORY_PAGE_TOO_LARGE`) (`createNormalizedMessage`/`generateMessageId`, 멀티 text 파트 id 충돌 방지 discriminator).
- `gjc-auth.provider.ts` — `command -v gjc` + 로그인 상태(agent.db:auth_credentials 존재 or `gjc` CLI). 미설치/미인증은 데이터로 반환(예외 아님).
- `gjc-skills.provider.ts` — `SkillsProvider` 확장. 루트: user `~/.gjc/agent/skills`, project `<ws>/.gjc/skills`. prefix: 스킬은 트리거 자동활성(명령형 아님) — codex `$`/claude `/` 참고해 gjc 표기 확정(잠정 `/`).
- `gjc-mcp.provider.ts` — gjc MCP 설정 위치 확정 필요(미조사). 최소 안전 stub(빈 목록) 또는 조사 후.
Expand Down
239 changes: 89 additions & 150 deletions server/modules/providers/list/gjc/gjc-sessions.provider.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,90 +5,15 @@ import type { IProviderSessions } from '@/shared/interfaces.js';
import type { AnyRecord, FetchHistoryOptions, FetchHistoryResult, NormalizedMessage } from '@/shared/types.js';
import { assignTranscriptTurns, type TranscriptTurnRecord } from '@/modules/providers/list/gjc/gjc-transcript-turns.js';
import { readGjcTranscriptMessage } from '@/modules/providers/list/gjc/gjc-transcript-message.js';
import { createNormalizedMessage, generateMessageId, readObjectRecord, sliceTailPage } from '@/shared/utils.js';
import { AppError, createNormalizedMessage, generateMessageId, readObjectRecord } from '@/shared/utils.js';

const PROVIDER = 'gjc';

const MAX_JSONL_LINE_BYTES = 32 * 1024 * 1024;
const MAX_BUFFERED_HISTORY_RECORDS = 5_000;
const MAX_BUFFERED_HISTORY_BYTES = 64 * 1024 * 1024;
const PAGINATION_RECORD_HEADROOM = 100;

type BufferedNormalizedMessage = {
message: NormalizedMessage;
byteLength: number;
};

/**
* Retains only the newest normalized transcript records. The byte limit accounts
* for the serialized record, which bounds the retained message strings and
* structured tool payloads without retaining an unbounded JSONL transcript.
*/
class NormalizedMessageRingBuffer {
private entries: Array<BufferedNormalizedMessage | undefined> = [];
private startIndex = 0;
private bufferedBytes = 0;

truncated = false;

constructor(
private readonly maxRecords: number,
private readonly maxBytes: number,
) {}

push(message: NormalizedMessage): void {
const byteLength = Buffer.byteLength(JSON.stringify(message), 'utf8');

if (byteLength > this.maxBytes) {
this.truncated = true;
return;
}

while (
this.entries.length - this.startIndex >= this.maxRecords
|| this.bufferedBytes + byteLength > this.maxBytes
) {
const oldest = this.entries[this.startIndex];
if (!oldest) {
break;
}
this.entries[this.startIndex] = undefined;
this.startIndex += 1;
this.bufferedBytes -= oldest.byteLength;
this.truncated = true;
}

this.entries.push({ message, byteLength });
this.bufferedBytes += byteLength;

if (this.startIndex >= 1_024) {
this.entries = this.entries.slice(this.startIndex);
this.startIndex = 0;
}
}

get messages(): NormalizedMessage[] {
const messages: NormalizedMessage[] = [];
for (let index = this.startIndex; index < this.entries.length; index += 1) {
const entry = this.entries[index];
if (entry) {
messages.push(entry.message);
}
}
return messages;
}
}

function getHistoryBufferRecordLimit(limit: number | null, offset: number): number {
if (limit === null) {
return MAX_BUFFERED_HISTORY_RECORDS;
}

return Math.min(
MAX_BUFFERED_HISTORY_RECORDS,
Math.max(PAGINATION_RECORD_HEADROOM, limit + offset + PAGINATION_RECORD_HEADROOM),
);
}
const HISTORY_READ_ATTEMPTS = 3;
type HistoryRow = { ordinal: number; time: number; kind: NormalizedMessage['kind']; toolId?: string };

/**
* Streams newline-delimited UTF-8 text while discarding a line as soon as it
Expand Down Expand Up @@ -197,20 +122,11 @@ async function readTranscriptLineage(sessionFilePath: string): Promise<Transcrip
* its own intermediate record with a unique id so multi-part turns never collide.
*/
async function streamGjcSessionMessages(
sessionId: string,
sessionFilePath: string,
turns: ReturnType<typeof assignTranscriptTurns>,
onMessage: (message: AnyRecord) => void,
): Promise<void> {
try {
const sessionFilePath = sessionsDb.getSessionById(sessionId)?.jsonl_path;

if (!sessionFilePath) {
console.warn(`gjc session file not found for session ${sessionId}`);
return;
}

// Which turn each record belongs to, and how that turn ended. Both come from
// the transcript, so a reloaded conversation reports what it reported live.
const turns = assignTranscriptTurns(await readTranscriptLineage(sessionFilePath));

for await (const line of readBoundedJsonlLines(sessionFilePath)) {
if (!line.trim()) {
Expand Down Expand Up @@ -354,12 +270,15 @@ async function streamGjcSessionMessages(
break;
}
}
} catch {
// Skip malformed lines.
} catch (error) {
if (error instanceof AppError) throw error;
// Skip malformed lines, not explicit page safety failures.
}
}
} catch (error) {
console.error(`Error reading gjc session messages for ${sessionId}:`, error);
if (error instanceof AppError) throw error;
console.error('Error reading gjc session messages:', error);
throw error;
}
}

Expand Down Expand Up @@ -494,72 +413,92 @@ export class GjcSessionsProvider implements IProviderSessions {
const { limit = null, offset = 0 } = options;
const normalizedOffset = Math.max(0, offset);
const normalizedLimit = limit === null ? null : Math.max(0, limit);
const messageBuffer = new NormalizedMessageRingBuffer(
getHistoryBufferRecordLimit(normalizedLimit, normalizedOffset),
MAX_BUFFERED_HISTORY_BYTES,
);

try {
await streamGjcSessionMessages(sessionId, (rawMessage) => {
for (const message of this.normalizeHistoryEntry(rawMessage, sessionId)) {
messageBuffer.push(message);
}
});
} catch (error) {
const message = error instanceof Error ? error.message : String(error);
console.warn(`[GjcProvider] Failed to load session ${sessionId}:`, message);
return {
messages: [],
total: 0,
hasMore: false,
offset: normalizedOffset,
limit: normalizedLimit,
};
const sessionFilePath = sessionsDb.getSessionById(sessionId)?.jsonl_path;
const empty = { messages: [], total: 0, hasMore: false, offset: normalizedOffset, limit: normalizedLimit };
if (!sessionFilePath) return empty;
// A live writer appends between the index and payload passes often enough
// that one changed revision is routine, not a client-visible failure.
for (let attempt = 1; ; attempt++) {
try {
return await this.readHistoryPage(sessionFilePath, sessionId, normalizedLimit, normalizedOffset);
} catch (error) {
if (error instanceof AppError && error.code === 'HISTORY_CHANGED' && attempt < HISTORY_READ_ATTEMPTS) continue;
if ((error as NodeJS.ErrnoException)?.code === 'ENOENT') return empty;
throw error;
}
}
}

const normalized = messageBuffer.messages.sort(
(a, b) => new Date(a.timestamp || 0).getTime() - new Date(b.timestamp || 0).getTime(),
);

const toolResultMap = new Map<string, NormalizedMessage>();
for (const msg of normalized) {
if (msg.kind === 'tool_result' && msg.toolId) {
toolResultMap.set(msg.toolId, msg);
private async readHistoryPage(
sessionFilePath: string,
sessionId: string,
normalizedLimit: number | null,
normalizedOffset: number,
): Promise<FetchHistoryResult> {
const revision = await fsSync.promises.stat(sessionFilePath);
const turns = assignTranscriptTurns(await readTranscriptLineage(sessionFilePath));
// Index only small descriptors. Tool-result rows must not consume visible
// pagination offsets or evict older messages from a payload ring.
const index: HistoryRow[] = [];
await streamGjcSessionMessages(sessionFilePath, turns, raw => {
for (const message of this.normalizeHistoryEntry(raw, sessionId)) {
const time = Date.parse(message.timestamp);
index.push({ ordinal: index.length, time: Number.isFinite(time) ? time : 0, kind: message.kind, toolId: message.toolId });
}
});
const chronological = index.sort((a, b) => a.time - b.time || a.ordinal - b.ordinal);
const visible = chronological.filter(row => row.kind !== 'tool_result');
const end = Math.max(0, visible.length - normalizedOffset);
if (normalizedLimit === null && end > MAX_BUFFERED_HISTORY_RECORDS) {
throw new AppError('History is too large to load at once; use paginated history.', { code: 'HISTORY_PAGE_TOO_LARGE', statusCode: 413 });
}
const start = Math.max(0, end - Math.min(normalizedLimit ?? MAX_BUFFERED_HISTORY_RECORDS, MAX_BUFFERED_HISTORY_RECORDS));
const selected = visible.slice(start, end);
const wanted = new Set(selected.map(row => row.ordinal));
const resultsByTool = new Map<string, number>();
const selectedTools = new Set(selected.filter(row => row.kind === 'tool_use' && row.toolId).map(row => row.toolId!));
for (const row of chronological) {
if (row.kind === 'tool_result' && row.toolId && selectedTools.has(row.toolId)) resultsByTool.set(row.toolId, row.ordinal);
}
for (const msg of normalized) {
if (msg.kind === 'tool_use' && msg.toolId && toolResultMap.has(msg.toolId)) {
const toolResult = toolResultMap.get(msg.toolId);
if (toolResult) {
// The standalone tool_result row is dropped below, so anything the
// UI needs has to be copied here. `toolUseResult` carries the
// runtime's typed details; omitting it silently made a reloaded
// transcript poorer than the turn that produced it.
msg.toolResult = {
content: toolResult.content,
isError: toolResult.isError,
...(toolResult.toolUseResult === undefined
? {}
: { toolUseResult: toolResult.toolUseResult }),
};
for (const ordinal of resultsByTool.values()) wanted.add(ordinal);

// A second streaming pass retains only the page and its attached results,
// independent of how far back the user has paged. No growing payload cache.
const payloads = new Map<number, NormalizedMessage>();
let ordinal = 0;
let bytes = 0;
if (wanted.size) await streamGjcSessionMessages(sessionFilePath, turns, raw => {
for (const message of this.normalizeHistoryEntry(raw, sessionId)) {
const position = ordinal++;
if (!wanted.has(position)) continue;
bytes += Buffer.byteLength(JSON.stringify(message), 'utf8');
if (bytes > MAX_BUFFERED_HISTORY_BYTES) {
throw new AppError('History page exceeds the payload limit; request fewer messages.', { code: 'HISTORY_PAGE_TOO_LARGE', statusCode: 413 });
}
payloads.set(position, message);
}
});
const after = await fsSync.promises.stat(sessionFilePath);
if (after.size !== revision.size || after.mtimeMs !== revision.mtimeMs) {
throw new AppError('Transcript changed while reading history; retry the request.', { code: 'HISTORY_CHANGED', statusCode: 409 });
}

// Tool results render inside their call, never as standalone timeline rows.
// When the bounded ring has discarded older rows, `total` is a lower bound;
// `hasMore` remains true so callers know the complete history was not retained.
const visibleMessages = normalized.filter((msg) => msg.kind !== 'tool_result');
const { page, hasMore: pageHasMore } = sliceTailPage(
visibleMessages,
normalizedLimit,
normalizedOffset,
);

const messages = selected.map(row => {
const message = payloads.get(row.ordinal)!;
const resultOrdinal = row.toolId ? resultsByTool.get(row.toolId) : undefined;
const result = resultOrdinal === undefined ? undefined : payloads.get(resultOrdinal);
if (message.kind === 'tool_use' && result) {
message.toolResult = {
content: result.content,
isError: result.isError,
...(result.toolUseResult === undefined ? {} : { toolUseResult: result.toolUseResult }),
};
}
return message;
});
return {
messages: page,
total: visibleMessages.length,
hasMore: pageHasMore || messageBuffer.truncated,
messages,
total: visible.length,
hasMore: start > 0,
offset: normalizedOffset,
limit: normalizedLimit,
tokenUsage: null,
Expand Down
8 changes: 3 additions & 5 deletions server/modules/providers/services/session-export.service.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
import path from 'node:path';

import { sessionsDb } from '@/modules/database/index.js';
import { providerRegistry } from '@/modules/providers/provider.registry.js';
import { fetchCompleteHistory } from '@/modules/providers/services/sessions.service.js';
import { sessionTranscriptWorkspace } from '@/modules/providers/services/session-worktrees.service.js';
import type { LLMProvider, NormalizedMessage } from '@/shared/types.js';
import { AppError } from '@/shared/utils.js';
Expand Down Expand Up @@ -192,10 +192,8 @@ export async function exportSessionTranscript(
let messages: NormalizedMessage[] = [];
const executionCwd = session.provider_session_id ? sessionTranscriptWorkspace(sessionId, projectPath) : null;
if (session.provider_session_id) {
const provider = providerRegistry.resolveProvider(session.provider as LLMProvider);
const history = await provider.sessions.fetchHistory(sessionId, {
limit: null,
offset: 0,
const history = await fetchCompleteHistory(sessionId, {
provider: session.provider as LLMProvider,
projectPath: executionCwd!,
providerSessionId: session.provider_session_id,
});
Expand Down
Loading