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
4 changes: 2 additions & 2 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -63,7 +63,7 @@ CHT_CORE_PATH=/path/to/cht-core
# Langfuse LLM observability (https://cloud.langfuse.com)
# Set LANGFUSE_ENABLED=false to disable tracing locally (default: enabled; the SDK
# disables itself when credentials are missing).
LANGFUSE_PUBLIC_KEY=pk-lf-...
LANGFUSE_SECRET_KEY=sk-lf-...
# LANGFUSE_PUBLIC_KEY=pk-lf-...
# LANGFUSE_SECRET_KEY=sk-lf-...
LANGFUSE_BASE_URL=https://cloud.langfuse.com
# LANGFUSE_ENABLED=false # uncomment to disable tracing
63 changes: 54 additions & 9 deletions docs/observability.md
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,37 @@ each time rather than mutating an earlier run's session. The PR identity lives i
All traces from a single `run-pipeline` invocation share one **session ID** (a UUID generated at the
start of `runPipeline`), visible in the Sessions view.

### Research and development CLIs

`npm run research`, `npm run dev:run`, and `npm run full` each produce **one trace per CLI run**
(`cht-agent-research` / `cht-agent-dev` / `cht-agent-full`). The supervisors add their own
observations under whatever observation is active, so a full run with human-feedback iterations shows
every research and development pass in one tree:

```
cht-agent-full (root span — input: { ticket })
├── agent: research (input: { title, domain, additionalContext? })
│ ├── retriever: documentation-search
│ ├── retriever: code-context-search
│ ├── retriever: context-analysis
│ └── chain: generate-plan
│ └── generation: orchestration-plan
└── agent: development (input: { title, previewMode, additionalContext? })
├── agent: generate-code (input: { iteration }; output: files, confidence, cross-file issues, model, tokens, cost)
│ └── generation: code-gen-plan, code-gen-execute, code-gen-execute-relaxed (claude-code-cli)
│ or code-gen-plan, code-gen-file, code-gen-continuation (claude-api)
├── evaluator: validate-implementation (output: { score })
│ └── generation: implementation-validation
└── … repeats per refinement iteration
```

Graph nodes report failures through the state's `errors` array rather than throwing, so a node span
is marked `level: 'ERROR'` whenever its update carries errors. Any ERROR observation (a failed node, a failed or
empty generation) also marks the **root** ERROR with that message, so failed runs are filterable at
the trace level even when the workflow recovered. The root output is a run summary (phase and errors
for research; approvals, iterations, score, and files written for dev/full). Node outputs are summaries (phase,
the node's status message, counts); full prompts and completions live only on the generations.

---

## Naming Conventions
Expand Down Expand Up @@ -125,7 +156,7 @@ Follow these when adding instrumentation to new workflows:
started before it are not retroactively updated.

3. Wrap each model call in `observeGeneration`. It records the prompt, the parsed output,
`usageDetails` (API path) or `costDetails` (Claude CLI path), and ends the generation with
`usageDetails`, `costDetails` when the path reports cost (Claude CLI), and ends the generation with
`level: 'ERROR'` on failure. Build the chain with `createLangChainStructuredChain` (API) or
`createStructuredCliChain` (CLI) so `invoke` returns a `GenerationResult`
(`{ parsed, model?, usage?, costUsd? }`):
Expand All @@ -134,26 +165,39 @@ Follow these when adding instrumentation to new workflows:
() => chain.invoke(prompt));
```

4. Use `root.startObservation()` for non-LangChain operations:
4. Code with no step handle (the code-gen modules) uses `observeActiveGeneration(opts, invoke)`,
which nests under the active observation and records nothing outside a trace. Both helpers take
`output` to map what is recorded and `failure` for a result that did not throw but still failed
(for example a CLI `is_error`).

5. For nested steps that don't have the root in scope, use `observeStep` (or `observeNode` for a
LangGraph node). Both nest under the **active** observation, so they need no parent argument, and
the callback receives the step to pass to `observeGeneration`:
```typescript
.addNode('generatePlan', observeNode({ name: 'generate-plan', asType: 'chain' }, this.generatePlanNode.bind(this)))
```
Until a `withTrace` has run in the process they are no-ops, so supervisor unit tests need no stubbing.

Use `root.startObservation()` for non-LangChain operations:
```typescript
const span = root.startObservation('my-operation', { input: { key: value } });
// ... do work ...
span.update({ output: { result } }).end();
```

5. Score terminal outcomes and set the root output:
6. Score terminal outcomes and set the root output:
```typescript
scoreTrace(root, { name: 'my-outcome', value: success ? 1 : 0 });
root.update({ output: { result } });
```

6. Shut Langfuse down **once** at the end of the whole run, not per item. It flushes buffered spans,
7. Shut Langfuse down **once** at the end of the whole run, not per item. It flushes buffered spans,
then queued scores, and awaits in-flight requests:
```typescript
await shutdownLangfuse();
```

7. In tests, stub the observability module with proxyquire. Record calls so the spec can assert the
8. In tests, stub the observability module with proxyquire. Record calls so the spec can assert the
root was actually passed to each stage (see `test/scripts/run-pipeline.spec.ts` for a full
recording spy):
```typescript
Expand Down Expand Up @@ -187,15 +231,16 @@ Every generation records the model, prompt, completion, latency, and errors. Cos
| Path | What is sent | How Langfuse prices it |
|---|---|---|
| API (OpenRouter / Anthropic) | `usageDetails` from LangChain `usage_metadata` (input/output/total tokens) and the provider's reported model name | Inferred from the model definition in Langfuse |
| Claude CLI (`claude -p`) | `costDetails.total` from the CLI's `total_cost_usd` and the configured model name; the CLI reports no token counts | Ingested USD directly |
| Claude CLI (`claude -p`) | `costDetails.total` from the CLI's `total_cost_usd`, `usageDetails` summed from its `modelUsage` (input includes cache reads/writes), and the model the CLI reports it ran (the highest-cost `modelUsage` entry), not the configured label | Ingested USD directly |

## Delivery Semantics

Spans are exported over OTLP/HTTP by a batching span processor with a 3-second request timeout. The
OTLP exporter applies its own retry/backoff on transient failures within that time. Scores go through
the client's queue, which uses the SDK default timeout of 60 seconds (`@langfuse/client` 5.11.1 ignores
its `timeout` option). `shutdownLangfuse()` flushes both queues before `process.exit` and only logs a
warning when Langfuse is unreachable, so tracing never changes a run's exit code.
its `timeout` option). `shutdownLangfuse()` flushes both queues before `process.exit`; if Langfuse is
unreachable it logs a `[Langfuse] … flush failed` warning and resolves, so tracing never changes a
run's exit code.

---

Expand All @@ -217,4 +262,4 @@ warning when Langfuse is unreachable, so tracing never changes a run's exit code
- **Dashboards**: cost per PR, filter decision distribution, distill success rate, latency by stage.
- **Alerting**: notify when `flag-for-human` rate or schema validation failures spike.
- **Evaluation datasets**: curate representative PR inputs from the trace store for automated quality evaluation.
- **Agent instrumentation**: apply the same pattern to `ResearchSupervisor` and individual agents once the memory pipeline PoC is validated.
- **Agent internals**: MCP/Kapa and DeepWiki calls inside the research retrievers are not yet separate `tool` observations.
11 changes: 11 additions & 0 deletions src/agents/code-generation-agent.ts
Original file line number Diff line number Diff line change
Expand Up @@ -355,10 +355,15 @@ export class CodeGenerationAgent {
crossFileIssues: crossFileIssues.length > 0 ? crossFileIssues : undefined,
compileGateSkipped: llmResult.compileGateSkipped,
compileGateSkipReason: llmResult.compileGateSkipReason,
tokensUsed: llmResult.tokensUsed,
modelUsed: llmResult.modelUsed,
costUsd: llmResult.costUsd,
};

console.log(`[Code Generation Agent] Generated ${result.files.length} files`);
console.log(`[Code Generation Agent] Confidence: ${(result.confidence * 100).toFixed(0)}%`);
const cost = result.costUsd === undefined ? 'n/a' : '$' + result.costUsd.toFixed(4);
console.log(`[Code Generation Agent] Model: ${result.modelUsed ?? 'unknown'}, tokens: ${result.tokensUsed ?? 'n/a'}, cost: ${cost}`);
this.todos.printSummary();
return result;
}
Expand Down Expand Up @@ -692,6 +697,9 @@ export class CodeGenerationAgent {
moduleCrossFileIssues?: import('../types').CrossFileIssue[];
compileGateSkipped?: boolean;
compileGateSkipReason?: string;
tokensUsed?: number;
modelUsed?: string;
costUsd?: number;
}> {
const moduleInput = this.buildModuleInput(input, context);
const session = await this.initBeadsSession(input);
Expand All @@ -711,6 +719,9 @@ export class CodeGenerationAgent {
moduleCrossFileIssues: moduleOutput.crossFileIssues,
compileGateSkipped: moduleOutput.compileGateSkipped,
compileGateSkipReason: moduleOutput.compileGateSkipReason,
tokensUsed: moduleOutput.tokensUsed,
modelUsed: moduleOutput.modelUsed,
costUsd: moduleOutput.costUsd,
};
}

Expand Down
18 changes: 15 additions & 3 deletions src/cli/dev.ts
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@

import * as dotenv from 'dotenv';
import * as path from 'node:path';
import { withTrace, shutdownLangfuse, type TraceRoot } from '../observability';
import { DevelopmentSupervisor } from '../supervisors/development-supervisor';
import { parseTicketFile } from '../utils/ticket-parser';
import { displayIssueDetails } from '../workflows/research-workflow';
Expand Down Expand Up @@ -167,7 +168,7 @@ function ensureTicketPath(): string {
return path.resolve(process.argv[2]);
}

const main = async (): Promise<void> => {
const main = async (root: TraceRoot): Promise<void> => {
console.log('╔════════════════════════════════════════════════════════════════╗');
console.log('║ CHT Multi-Agent System - Development Only CLI ║');
console.log('║ (Research phase skipped — using synthesized stubs) ║');
Expand Down Expand Up @@ -217,16 +218,27 @@ const main = async (): Promise<void> => {

// Display completion
displayDevelopmentCompletion(workflowResult, developmentInput.options);
root.update({ output: {
approved: workflowResult.approved,
iterations: workflowResult.iterationCount,
score: workflowResult.result?.validationResult?.overallScore,
filesWritten: workflowResult.filesWritten,
} });

} catch (error) {
console.error('\n❌ Error running development:', error);
if (error instanceof Error) {
console.error('Message:', error.message);
console.error('Stack:', error.stack);
}
process.exit(1);
throw error;
}
};

// Run the CLI
main();
withTrace({ name: 'cht-agent-dev', tags: ['cht-agent', 'dev'], input: { ticket: process.argv[2] } }, main)
.then(shutdownLangfuse)
.catch(async () => {
await shutdownLangfuse();
process.exit(1);
});
1 change: 1 addition & 0 deletions src/cli/display-helpers.ts
Original file line number Diff line number Diff line change
Expand Up @@ -296,6 +296,7 @@ export const runResearchWorkflow = async (
const duration = ((Date.now() - startTime) / 1000).toFixed(2);

displayResults(result, duration);
return result;
};

export const displayResults = (result: ResearchState, duration: string) => {
Expand Down
19 changes: 16 additions & 3 deletions src/cli/full.ts
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@

import * as dotenv from 'dotenv';
import * as path from 'node:path';
import { withTrace, shutdownLangfuse, type TraceRoot } from '../observability';
import { ResearchSupervisor } from '../supervisors/research-supervisor';
import { DevelopmentSupervisor } from '../supervisors/development-supervisor';
import { parseTicketFile } from '../utils/ticket-parser';
Expand Down Expand Up @@ -86,7 +87,7 @@ function ensureTicketPath(): string {
return path.resolve(process.argv[2]);
}

const main = async (): Promise<void> => {
const main = async (root: TraceRoot): Promise<void> => {
console.log('╔════════════════════════════════════════════════════════════════╗');
console.log('║ CHT Multi-Agent System - Full Workflow CLI ║');
console.log('╚════════════════════════════════════════════════════════════════╝\n');
Expand Down Expand Up @@ -117,15 +118,27 @@ const main = async (): Promise<void> => {
developmentOptions
);
displayFullWorkflowSummary(workflowResult);
root.update({ output: {
researchApproved: workflowResult.research.approved,
researchIterations: workflowResult.research.iterationCount,
developmentApproved: workflowResult.development?.approved,
developmentIterations: workflowResult.development?.iterationCount,
filesWritten: workflowResult.development?.filesWritten,
} });
} catch (error) {
console.error('\n❌ Error running workflow:', error);
if (error instanceof Error) {
console.error('Message:', error.message);
console.error('Stack:', error.stack);
}
process.exit(1);
throw error;
}
};

// Run the CLI
main();
withTrace({ name: 'cht-agent-full', tags: ['cht-agent', 'full'], input: { ticket: process.argv[2] } }, main)
.then(shutdownLangfuse)
.catch(async () => {
await shutdownLangfuse();
process.exit(1);
});
11 changes: 6 additions & 5 deletions src/cli/research.ts
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@
import * as dotenv from 'dotenv';
import * as path from 'node:path';
import { runResearchWorkflow } from './display-helpers';
import { withTrace, shutdownLangfuse } from '../observability';

dotenv.config();

Expand All @@ -37,15 +38,15 @@ const HELP_HINTS = [
'Content goes in markdown body with ## sections',
];

runResearchWorkflow(
'CHT Multi-Agent System - Research CLI',
getTicketPath,
HELP_HINTS,
).catch(error => {
withTrace({ name: 'cht-agent-research', tags: ['cht-agent', 'research'], input: { ticket: process.argv[2] } }, async root => {
const result = await runResearchWorkflow('CHT Multi-Agent System - Research CLI', getTicketPath, HELP_HINTS);
root.update({ output: { phase: result.currentPhase, planGenerated: result.orchestrationPlan !== undefined, errors: result.errors } });
}).then(shutdownLangfuse).catch(async error => {
console.error('\n❌ Error running research workflow:', error);
if (error instanceof Error) {
console.error('Message:', error.message);
console.error('Stack:', error.stack);
}
await shutdownLangfuse();
process.exit(1);
});
1 change: 1 addition & 0 deletions src/layers/code-gen/interface.ts
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,7 @@ export interface CodeGenModuleOutput {
explanation: string;
tokensUsed?: number;
modelUsed?: string;
costUsd?: number;
/**
* True when the module knows its output is incomplete (e.g., CLI hit
* is_error or saturated max-turns). The agent surfaces this as a
Expand Down
18 changes: 13 additions & 5 deletions src/layers/code-gen/modules/claude-api/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,9 +5,10 @@ import {
ContextFile,
GeneratedFile,
} from '../../interface';
import { LLMProvider, LLMToolDefinition, ToolHandler, createAnthropicProvider, getAPIConfigFromEnv } from '../../../../llm';
import { InvokeOptions, LLMProvider, LLMResponse, LLMToolDefinition, ToolHandler, createAnthropicProvider, getAPIConfigFromEnv } from '../../../../llm';
import { readEnv } from '../../../../utils/env';
import { isShutdownRequested } from '../../../../utils/shutdown';
import { fromLLMResponse, observeActiveGeneration } from '../../../../observability';
import { runApiCompileGate } from './compile-gate';
import { CompileValidationResult } from '../../../../agents/compile-validator';
import {
Expand Down Expand Up @@ -403,10 +404,9 @@ export class ClaudeApiCodeGenModule implements CodeGenModule {
input: CodeGenModuleInput,
manifest: FileManifest
): Promise<{ plan: PlanItem[]; tokensUsed: number }> {
const llm = this.getProvider();
const prompt = this.buildPlanPrompt(input, manifest);

const response = await llm.invoke(prompt, { maxTokens: 8192 });
const response = await this.tracedInvoke('code-gen-plan', prompt, { maxTokens: 8192 });
const tokensUsed = (response.usage?.inputTokens ?? 0) + (response.usage?.outputTokens ?? 0);

const plan = this.parsePlan(response.content);
Expand Down Expand Up @@ -678,13 +678,21 @@ export class ClaudeApiCodeGenModule implements CodeGenModule {
};
}

private tracedInvoke(name: string, prompt: string, options: InvokeOptions): Promise<LLMResponse> {
const llm = this.getProvider();
return observeActiveGeneration({ name, model: llm.modelName, input: prompt, output: (r) => r.content }, async () => {
const response = await llm.invoke(prompt, options);
return fromLLMResponse(response, response);
});
}

private async invokeLLM(
prompt: string,
codeGenTools: { tools: LLMToolDefinition[]; toolHandler: ToolHandler } | undefined,
filePath: string,
): Promise<Awaited<ReturnType<LLMProvider['invoke']>> | null> {
try {
return await this.getProvider().invoke(prompt, {
return await this.tracedInvoke('code-gen-file', prompt, {
maxTokens: 65536,
...(codeGenTools ? { tools: codeGenTools.tools, toolHandler: codeGenTools.toolHandler } : {}),
});
Expand Down Expand Up @@ -838,7 +846,7 @@ export class ClaudeApiCodeGenModule implements CodeGenModule {
const prompt = this.buildContinuationPrompt(lastLines, planItem, input);
let response;
try {
response = await this.getProvider().invoke(prompt, { maxTokens: 65536 });
response = await this.tracedInvoke('code-gen-continuation', prompt, { maxTokens: 65536 });
} catch (error) {
console.error(`[Code Gen Module] Continuation call ${iteration + 1} failed:`, error);
return false;
Expand Down
Loading