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
58 changes: 46 additions & 12 deletions docs/proposals/110-usage-metrics-local-runs.md
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@

**Author:** Bubba (VoynichLabs)
**Date:** 2026-03-08
**Status:** Proposal
**Status:** Implemented (PR #219)

---

Expand Down Expand Up @@ -165,18 +165,52 @@ None of these are present in a bare CLI run. The file-based approach requires on

---

## Open Questions
## Implementation Notes (PR #219)

1. **Ollama field mapping**: `prompt_eval_count` / `eval_count` vs `prompt_tokens` /
`completion_tokens` — should `extract_token_count()` add an Ollama-specific branch, or
should the Ollama LLM adapter normalize field names before returning?
The final implementation diverges from Option A above in several ways:

2. **Thread safety**: Luigi runs tasks in parallel workers. The appender needs an
`fcntl.flock()` or atomic rename approach to avoid corruption under concurrent writes.
### File format: JSONL, not JSON

3. **Thinking tokens for local models**: Qwen 3 thinking tokens are stripped before the
response reaches llama_index (LM Studio does not expose them separately). The
`thinking_tokens` field will be `null` for local runs — acceptable for now.
`usage_metrics.jsonl` uses one JSON object per line (append-only). This avoids the need for
atomic rename or in-memory accumulation — each LLM call appends a single line. Thread-safe
by nature since each write is a short append to a file handle.

4. **Retention**: Should `usage_metrics.json` be overwritten on resume (Luigi skip-completed
behavior) or appended? Appending is safer — a resumed run adds only the new calls.
### Recording source: llama_index instrumentation, not LLMExecutor

Successful calls are recorded by `TrackActivity` (the llama_index `BaseEventHandler`) which
receives the actual `ChatResponse` with full token counts, cost, and `provider:model` info.
`LLMExecutor._record_attempt_token_metrics()` only records **failures**, since instrumentation
end events are not emitted when the LLM call fails.

This was necessary because `execute_function(llm)` returns the processed result (a Pydantic
model or string), not the raw `ChatResponse`. The instrumentation layer is the only place
with access to the real response.

### Model field includes provider

The `model` field contains the full `provider:model` string (e.g.
`Google AI Studio:google/gemini-2.0-flash-001`), matching `activity_overview.json`.

### Example output

```json
{"timestamp": "2026-03-10T13:36:48.250446", "success": true, "model": "Google AI Studio:google/gemini-2.0-flash-001", "duration_seconds": 4.879, "input_tokens": 5316, "output_tokens": 643, "cost_usd": 0.0007888}
{"timestamp": "2026-03-10T13:36:53.554864", "success": true, "model": "Google:google/gemini-2.0-flash-001", "duration_seconds": 5.237, "input_tokens": 8877, "output_tokens": 562, "cost_usd": 0.0011125}
```

### Key files

| File | Role |
|------|------|
| `worker_plan/worker_plan_internal/llm_util/usage_metrics.py` | Core module: `set_usage_metrics_path()`, `record_usage_metric()` |
| `worker_plan/worker_plan_internal/llm_util/track_activity.py` | Records successful calls via `_record_file_usage_metric()` |
| `worker_plan/worker_plan_internal/llm_util/llm_executor.py` | Records failed calls only |
| `worker_plan/worker_plan_internal/plan/run_plan_pipeline.py` | Sets/clears metrics path around pipeline execution |
| `worker_plan/worker_plan_api/filenames.py` | `USAGE_METRICS_JSONL` constant |

### Resolved open questions

1. **Ollama field mapping**: Handled by `extract_token_count()` and `TrackActivity._extract_token_usage()` which already support multiple field name variations.
2. **Thread safety**: JSONL append-per-line is safe for concurrent Luigi workers.
3. **Thinking tokens**: Recorded when available (e.g. from OpenRouter reasoning models). `null` for providers that don't expose them.
4. **Retention on resume**: Appended. A resumed run adds only the new calls alongside the restored snapshot's existing metrics.
4 changes: 2 additions & 2 deletions docs/proposals/111-promising-directions.md
Original file line number Diff line number Diff line change
Expand Up @@ -34,7 +34,7 @@ Agents need to discover PlanExe, understand its tools, and consume outputs progr
|---|----------|-------------|
| **86** | Agent-Optimized Pipeline | Removes the 5 key friction points for autonomous agent use: human approval gate, no agent prompt examples, poll intervals tuned for humans, no machine-readable output, no autonomous agent setup docs |
| **62** | Agent-First Frontend Discoverability | `llms.txt`, `/.well-known/mcp.json`, agent-readable README — standard discovery protocols so agents find PlanExe without human guidance |
| **110** | Usage Metrics for Local Runs | Agents need cost accounting for budget-constrained workflows. Answers "how much did this run cost?" |
| **110** | Usage Metrics for Local Runs | ✅ **Implemented (PR #219)**. Agents need cost accounting for budget-constrained workflows. `usage_metrics.jsonl` answers "how much did this run cost?" with per-call granularity (model, tokens, cost, duration). Complements `activity_overview.json` aggregated totals |

Key friction points from #86 that block autonomous agent use:
- **F1**: Human approval step before `plan_create` — autonomous agents can't proceed
Expand Down Expand Up @@ -94,7 +94,7 @@ Phase 1: Reliable foundation (now)
├─ #87 Plan resume
├─ #109 Retry improvements
├─ #102 Error-feedback retries
├─ #110 Usage metrics
├─ #110 Usage metrics ✅
└─ #58 Prompt boost

Phase 2: Agent-native interface (next)
Expand Down
6 changes: 5 additions & 1 deletion frontend_multi_user/src/app.py
Original file line number Diff line number Diff line change
Expand Up @@ -3405,9 +3405,13 @@ def plan_resume():
return redirect(url_for('plan', id=run_id))

# Reject resume if the snapshot was created by a different pipeline version.
# Plans created before pipeline_version was stamped have no stored version;
# allow those through — the worker-side check in worker_plan_database/app.py
# reads pipeline_version from the actual snapshot metadata file
# (001-3-planexe_metadata.json) and rejects incompatible versions there.
stored_params = task.parameters if isinstance(task.parameters, dict) else {}
stored_version = stored_params.get("pipeline_version")
if stored_version != PIPELINE_VERSION:
if stored_version is not None and stored_version != PIPELINE_VERSION:
return redirect(url_for(
'plan', id=run_id,
resume_error="version_mismatch",
Expand Down
6 changes: 5 additions & 1 deletion mcp_cloud/db_queries.py
Original file line number Diff line number Diff line change
Expand Up @@ -268,9 +268,13 @@ def _resume_plan_sync(plan_id: str, model_profile: str) -> Optional[dict[str, An
}

# Reject resume if the snapshot was created by a different pipeline version.
# Plans created before pipeline_version was stamped have no stored version;
# allow those through — the worker-side check in worker_plan_database/app.py
# reads pipeline_version from the actual snapshot metadata file
# (001-3-planexe_metadata.json) and rejects incompatible versions there.
stored_params = plan.parameters if isinstance(plan.parameters, dict) else {}
stored_version = stored_params.get("pipeline_version")
if stored_version != PIPELINE_VERSION:
if stored_version is not None and stored_version != PIPELINE_VERSION:
return {
"error": {
"code": "PIPELINE_VERSION_MISMATCH",
Expand Down
1 change: 1 addition & 0 deletions worker_plan/worker_plan_api/filenames.py
Original file line number Diff line number Diff line change
Expand Up @@ -127,3 +127,4 @@ class ExtraFilenameEnum(str, Enum):
PIPELINE_STOP_REQUESTED_FLAG = "pipeline_stop_requested.txt"
TRACK_ACTIVITY_JSONL = "track_activity.jsonl"
ACTIVITY_OVERVIEW_JSON = "activity_overview.json"
USAGE_METRICS_JSONL = "usage_metrics.jsonl"
15 changes: 15 additions & 0 deletions worker_plan/worker_plan_internal/llm_util/llm_executor.py
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,7 @@
from llama_index.core.llms.llm import LLM
from llama_index.core.instrumentation.dispatcher import instrument_tags
from worker_plan_internal.llm_factory import get_llm
from worker_plan_internal.llm_util.usage_metrics import record_usage_metric

logger = logging.getLogger(__name__)

Expand Down Expand Up @@ -271,6 +272,20 @@ def _record_attempt_token_metrics(
except Exception as exc:
logger.debug("Failed to record token metrics for attempt: %s", exc)

# File-based usage metrics for local runs (no database required).
# Successful calls are recorded by TrackActivity via llama_index
# instrumentation, which has access to the real ChatResponse with
# full token counts and upstream provider/model info.
# Here we only record failures, since instrumentation end events
# are not emitted when the LLM call fails.
if not success:
record_usage_metric(
model=llm_model_name,
duration_seconds=duration,
success=False,
error_message=error_message,
)

def _check_stop_callback(self, last_attempt: LLMAttempt, start_time: float, attempt_index: int) -> None:
"""Checks the callback, if it exists, to see if execution should stop."""
if self.should_stop_callback is None:
Expand Down
18 changes: 18 additions & 0 deletions worker_plan/worker_plan_internal/llm_util/track_activity.py
Original file line number Diff line number Diff line change
Expand Up @@ -299,6 +299,23 @@ def _record_token_metrics_row(self, event_data: dict, duration_seconds: Optional
except Exception:
logger.debug("Failed to persist token metrics from TrackActivity", exc_info=True)

def _record_file_usage_metric(self, event_data: dict, token_usage: Optional[dict], duration_seconds: Optional[float]) -> None:
"""Write a usage metric row to usage_metrics.jsonl via the shared module."""
cost = self._extract_cost(event_data)
if not token_usage and cost == 0.0:
return
from worker_plan_internal.llm_util.usage_metrics import record_usage_metric
model_name = self._extract_model_name(event_data)
record_usage_metric(
model=model_name or "unknown",
duration_seconds=duration_seconds or 0.0,
success=True,
input_tokens=int(token_usage["input_tokens"]) if token_usage and token_usage.get("input_tokens") is not None else None,
output_tokens=int(token_usage["output_tokens"]) if token_usage and token_usage.get("output_tokens") is not None else None,
thinking_tokens=int(token_usage["thinking_tokens"]) if token_usage and token_usage.get("thinking_tokens") is not None else None,
cost_usd=cost or None,
)

@staticmethod
def _split_provider_and_model(model_name: str) -> tuple[Optional[str], Optional[str]]:
if not model_name:
Expand Down Expand Up @@ -367,6 +384,7 @@ def handle(self, event: Any) -> None:
event_record["token_usage"] = token_usage
self._update_activity_overview(filtered_event_data)
self._record_token_metrics_row(filtered_event_data, duration_seconds=duration_seconds)
self._record_file_usage_metric(filtered_event_data, token_usage, duration_seconds)
elif isinstance(event, (LLMChatStartEvent, LLMCompletionStartEvent, LLMStructuredPredictStartEvent)):
match_key = self._event_match_key(filtered_event_data)
start_ts = self._parse_event_timestamp(filtered_event_data)
Expand Down
81 changes: 81 additions & 0 deletions worker_plan/worker_plan_internal/llm_util/usage_metrics.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,81 @@
"""
File-based usage metrics for local runs.

Records per-LLM-call metrics (model, tokens, duration, success/failure) to a JSONL file
in the run output directory. Works without a database — designed for local/offline runs.

Usage:
from worker_plan_internal.llm_util.usage_metrics import set_usage_metrics_path, record_usage_metric

# Set once at pipeline start
set_usage_metrics_path(run_id_dir / ExtraFilenameEnum.USAGE_METRICS_JSONL.value)

# Called automatically by LLMExecutor._record_attempt_token_metrics()
record_usage_metric(model="gpt-4", duration=1.23, success=True, input_tokens=100, output_tokens=50)

# Clear after pipeline completes to avoid stale state
set_usage_metrics_path(None)
"""
import json
import logging
from datetime import datetime
from pathlib import Path
from typing import Optional

logger = logging.getLogger(__name__)

_usage_metrics_path: Optional[Path] = None


def set_usage_metrics_path(path: Optional[Path]) -> None:
"""Set the JSONL file path for recording usage metrics."""
global _usage_metrics_path
_usage_metrics_path = path


def get_usage_metrics_path() -> Optional[Path]:
"""Get the current JSONL file path for recording usage metrics."""
return _usage_metrics_path


def record_usage_metric(
model: str,
duration_seconds: float,
success: bool,
error_message: Optional[str] = None,
input_tokens: Optional[int] = None,
output_tokens: Optional[int] = None,
thinking_tokens: Optional[int] = None,
cost_usd: Optional[float] = None,
) -> None:
"""Append a single usage metric record to the JSONL file.

Best-effort: never raises exceptions to avoid blocking the LLM pipeline.
"""
path = _usage_metrics_path
if path is None:
logger.warning("record_usage_metric called but no usage metrics path is set")
return

record = {
"timestamp": datetime.now().isoformat(),
"success": success,
"model": model,
"duration_seconds": round(duration_seconds, 3),
}
if error_message:
record["error"] = error_message
if input_tokens is not None:
record["input_tokens"] = input_tokens
if output_tokens is not None:
record["output_tokens"] = output_tokens
if thinking_tokens is not None:
record["thinking_tokens"] = thinking_tokens
if cost_usd is not None:
record["cost_usd"] = cost_usd

try:
with open(path, "a", encoding="utf-8") as f:
f.write(json.dumps(record) + "\n")
except Exception as exc:
logger.warning("Failed to write usage metric: %s", exc)
10 changes: 10 additions & 0 deletions worker_plan/worker_plan_internal/plan/run_plan_pipeline.py
Original file line number Diff line number Diff line change
Expand Up @@ -3973,6 +3973,7 @@ def get_progress_percentage(self) -> PipelineProgress:
ExtraFilenameEnum.EXPECTED_FILENAMES1_JSON.value,
ExtraFilenameEnum.LOG_TXT.value,
ExtraFilenameEnum.PIPELINE_STOP_REQUESTED_FLAG.value,
ExtraFilenameEnum.USAGE_METRICS_JSONL.value,
'.DS_Store',
]
files = [f for f in files if f not in ignore_files]
Expand Down Expand Up @@ -4054,6 +4055,12 @@ def run(self):
logger.debug(f"Removing pre-existing stop flag file: {stop_flag_path!r}")
stop_flag_path.unlink()

# Enable file-based usage metrics for this run.
from worker_plan_internal.llm_util.usage_metrics import set_usage_metrics_path
usage_metrics_path = self.run_id_dir / ExtraFilenameEnum.USAGE_METRICS_JSONL.value
set_usage_metrics_path(usage_metrics_path)
logger.info(f"Usage metrics will be written to {usage_metrics_path}")

# create a json file with the expected filenames. Save it to the run/run_id/expected_filenames1.json
expected_filenames_path = self.run_id_dir / ExtraFilenameEnum.EXPECTED_FILENAMES1_JSON.value
with open(expected_filenames_path, "w") as f:
Expand All @@ -4074,6 +4081,9 @@ def run(self):
workers=luigi_workers
)

# Clear the usage metrics path after the run.
set_usage_metrics_path(None)

# After the pipeline finishes (or fails), check for the stop flag.
if self.has_stop_flag_file:
logger.info("Pipeline was stopped intentionally via PipelineStopRequested exception.")
Expand Down
9 changes: 8 additions & 1 deletion worker_plan_database/app.py
Original file line number Diff line number Diff line change
Expand Up @@ -442,7 +442,14 @@ def _handle_task_completion(self, parameters: HandleTaskCompletionParameters) ->
"""
logger.debug(f"ServerExecutePipeline._handle_task_completion")

WorkerItem.upsert_heartbeat(worker_id=WORKER_ID, current_task_id=self.task_id)
try:
WorkerItem.upsert_heartbeat(worker_id=WORKER_ID, current_task_id=self.task_id)
except Exception as exc:
logger.warning("Heartbeat upsert failed (non-fatal): %s", exc)
try:
db.session.rollback()
except Exception:
pass

# Lookup the taskitem in the database by self.task_id
task = db.session.get(PlanItem, self.task_id)
Expand Down