|
@@ -159,16 +159,22 @@ class SQLiteSessionStore:
|
|
|
self._touch_session_locked(session_id, now)
|
|
self._touch_session_locked(session_id, now)
|
|
|
self._connect().commit()
|
|
self._connect().commit()
|
|
|
|
|
|
|
|
- def complete_turn(self, session_id: str, *, turn_index: int) -> None:
|
|
|
|
|
|
|
+ def complete_turn(
|
|
|
|
|
+ self,
|
|
|
|
|
+ session_id: str,
|
|
|
|
|
+ *,
|
|
|
|
|
+ turn_index: int,
|
|
|
|
|
+ wall_time_ms: int | None = None,
|
|
|
|
|
+ ) -> None:
|
|
|
now = self._now()
|
|
now = self._now()
|
|
|
with self._lock:
|
|
with self._lock:
|
|
|
self._connect().execute(
|
|
self._connect().execute(
|
|
|
"""
|
|
"""
|
|
|
UPDATE turns
|
|
UPDATE turns
|
|
|
- SET completed_at = ?
|
|
|
|
|
|
|
+ SET completed_at = ?, wall_time_ms = ?
|
|
|
WHERE session_id = ? AND turn_index = ?
|
|
WHERE session_id = ? AND turn_index = ?
|
|
|
""",
|
|
""",
|
|
|
- (now, session_id, turn_index),
|
|
|
|
|
|
|
+ (now, wall_time_ms, session_id, turn_index),
|
|
|
)
|
|
)
|
|
|
self._touch_session_locked(session_id, now)
|
|
self._touch_session_locked(session_id, now)
|
|
|
self._connect().commit()
|
|
self._connect().commit()
|
|
@@ -274,29 +280,42 @@ class SQLiteSessionStore:
|
|
|
usage: TokenUsage,
|
|
usage: TokenUsage,
|
|
|
ttft_ms: int | None,
|
|
ttft_ms: int | None,
|
|
|
elapsed_ms: int,
|
|
elapsed_ms: int,
|
|
|
|
|
+ metadata: dict[str, Any] | None = None,
|
|
|
) -> None:
|
|
) -> None:
|
|
|
|
|
+ metadata = metadata or {}
|
|
|
|
|
+ call_kind = metadata.get("call_kind", "chat_completion")
|
|
|
|
|
+ persisted_usage = TokenUsage() if call_kind == "tool_execution" else usage
|
|
|
now = self._now()
|
|
now = self._now()
|
|
|
with self._lock:
|
|
with self._lock:
|
|
|
self._connect().execute(
|
|
self._connect().execute(
|
|
|
"""
|
|
"""
|
|
|
- INSERT INTO usage_stats
|
|
|
|
|
|
|
+ INSERT OR IGNORE INTO usage_stats
|
|
|
(
|
|
(
|
|
|
session_id, turn_index, round_index, prompt_tokens,
|
|
session_id, turn_index, round_index, prompt_tokens,
|
|
|
completion_tokens, total_tokens, cached_tokens,
|
|
completion_tokens, total_tokens, cached_tokens,
|
|
|
- ttft_ms, elapsed_ms, created_at
|
|
|
|
|
|
|
+ ttft_ms, elapsed_ms, agent, call_kind, event_id,
|
|
|
|
|
+ event_name, used_fallback, tool_latency_ms,
|
|
|
|
|
+ metric_key, created_at
|
|
|
)
|
|
)
|
|
|
- VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
|
|
|
|
|
|
+ VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
|
|
|
""",
|
|
""",
|
|
|
(
|
|
(
|
|
|
session_id,
|
|
session_id,
|
|
|
turn_index,
|
|
turn_index,
|
|
|
round_index,
|
|
round_index,
|
|
|
- usage.prompt_tokens,
|
|
|
|
|
- usage.completion_tokens,
|
|
|
|
|
- usage.total_tokens,
|
|
|
|
|
- usage.cached_tokens,
|
|
|
|
|
|
|
+ persisted_usage.prompt_tokens,
|
|
|
|
|
+ persisted_usage.completion_tokens,
|
|
|
|
|
+ persisted_usage.total_tokens,
|
|
|
|
|
+ persisted_usage.cached_tokens,
|
|
|
ttft_ms,
|
|
ttft_ms,
|
|
|
elapsed_ms,
|
|
elapsed_ms,
|
|
|
|
|
+ metadata.get("agent", "chat_agent"),
|
|
|
|
|
+ call_kind,
|
|
|
|
|
+ metadata.get("event_id"),
|
|
|
|
|
+ metadata.get("event_name"),
|
|
|
|
|
+ bool(metadata.get("used_fallback", False)),
|
|
|
|
|
+ metadata.get("tool_latency_ms"),
|
|
|
|
|
+ metadata.get("metric_key"),
|
|
|
now,
|
|
now,
|
|
|
),
|
|
),
|
|
|
)
|
|
)
|
|
@@ -352,31 +371,52 @@ class SQLiteSessionStore:
|
|
|
]
|
|
]
|
|
|
|
|
|
|
|
def usage_summary(self, session_id: str) -> dict[str, Any]:
|
|
def usage_summary(self, session_id: str) -> dict[str, Any]:
|
|
|
|
|
+ session_record = self.get_session(session_id)
|
|
|
|
|
+ mode = (
|
|
|
|
|
+ session_record["config"].get("tool_invocation_mode", "dual_agent")
|
|
|
|
|
+ if session_record is not None
|
|
|
|
|
+ else "dual_agent"
|
|
|
|
|
+ )
|
|
|
call_rows = self._fetchall(
|
|
call_rows = self._fetchall(
|
|
|
"""
|
|
"""
|
|
|
SELECT
|
|
SELECT
|
|
|
id, turn_index, round_index, prompt_tokens, completion_tokens,
|
|
id, turn_index, round_index, prompt_tokens, completion_tokens,
|
|
|
- total_tokens, cached_tokens, ttft_ms, elapsed_ms, created_at
|
|
|
|
|
|
|
+ total_tokens, cached_tokens, ttft_ms, elapsed_ms, agent,
|
|
|
|
|
+ call_kind, event_id, event_name, used_fallback,
|
|
|
|
|
+ tool_latency_ms, metric_key, created_at
|
|
|
FROM usage_stats
|
|
FROM usage_stats
|
|
|
WHERE session_id = ?
|
|
WHERE session_id = ?
|
|
|
ORDER BY id ASC
|
|
ORDER BY id ASC
|
|
|
""",
|
|
""",
|
|
|
(session_id,),
|
|
(session_id,),
|
|
|
)
|
|
)
|
|
|
- calls = [self._usage_row(row) for row in call_rows]
|
|
|
|
|
|
|
+ calls = [self._usage_row(row, mode=mode) for row in call_rows]
|
|
|
turn_rows = self._fetchall(
|
|
turn_rows = self._fetchall(
|
|
|
"""
|
|
"""
|
|
|
SELECT
|
|
SELECT
|
|
|
- turn_index,
|
|
|
|
|
- COALESCE(SUM(prompt_tokens), 0) AS prompt_tokens,
|
|
|
|
|
- COALESCE(SUM(completion_tokens), 0) AS completion_tokens,
|
|
|
|
|
- COALESCE(SUM(total_tokens), 0) AS total_tokens,
|
|
|
|
|
- COALESCE(SUM(cached_tokens), 0) AS cached_tokens,
|
|
|
|
|
- COALESCE(SUM(elapsed_ms), 0) AS elapsed_ms
|
|
|
|
|
- FROM usage_stats
|
|
|
|
|
- WHERE session_id = ?
|
|
|
|
|
- GROUP BY turn_index
|
|
|
|
|
- ORDER BY turn_index ASC
|
|
|
|
|
|
|
+ turns.turn_index,
|
|
|
|
|
+ COALESCE(SUM(CASE WHEN usage_stats.call_kind != 'tool_execution'
|
|
|
|
|
+ THEN usage_stats.prompt_tokens ELSE 0 END), 0) AS prompt_tokens,
|
|
|
|
|
+ COALESCE(SUM(CASE WHEN usage_stats.call_kind != 'tool_execution'
|
|
|
|
|
+ THEN usage_stats.completion_tokens ELSE 0 END), 0) AS completion_tokens,
|
|
|
|
|
+ COALESCE(SUM(CASE WHEN usage_stats.call_kind != 'tool_execution'
|
|
|
|
|
+ THEN usage_stats.total_tokens ELSE 0 END), 0) AS total_tokens,
|
|
|
|
|
+ COALESCE(SUM(CASE WHEN usage_stats.call_kind != 'tool_execution'
|
|
|
|
|
+ THEN usage_stats.cached_tokens ELSE 0 END), 0) AS cached_tokens,
|
|
|
|
|
+ COALESCE(SUM(usage_stats.elapsed_ms), 0) AS elapsed_ms,
|
|
|
|
|
+ COUNT(usage_stats.id) AS call_count,
|
|
|
|
|
+ COALESCE(SUM(CASE WHEN usage_stats.call_kind = 'argument_fallback'
|
|
|
|
|
+ THEN 1 ELSE 0 END), 0) AS fallback_count,
|
|
|
|
|
+ COALESCE(SUM(CASE WHEN usage_stats.call_kind = 'tool_execution'
|
|
|
|
|
+ THEN 1 ELSE 0 END), 0) AS tool_count,
|
|
|
|
|
+ turns.wall_time_ms AS turn_wall_time_ms
|
|
|
|
|
+ FROM turns
|
|
|
|
|
+ JOIN usage_stats
|
|
|
|
|
+ ON usage_stats.session_id = turns.session_id
|
|
|
|
|
+ AND usage_stats.turn_index = turns.turn_index
|
|
|
|
|
+ WHERE turns.session_id = ?
|
|
|
|
|
+ GROUP BY turns.turn_index, turns.wall_time_ms
|
|
|
|
|
+ ORDER BY turns.turn_index ASC
|
|
|
""",
|
|
""",
|
|
|
(session_id,),
|
|
(session_id,),
|
|
|
)
|
|
)
|
|
@@ -388,22 +428,43 @@ class SQLiteSessionStore:
|
|
|
"total_tokens": row["total_tokens"],
|
|
"total_tokens": row["total_tokens"],
|
|
|
"cached_tokens": row["cached_tokens"],
|
|
"cached_tokens": row["cached_tokens"],
|
|
|
"elapsed_ms": row["elapsed_ms"],
|
|
"elapsed_ms": row["elapsed_ms"],
|
|
|
|
|
+ "call_count": row["call_count"],
|
|
|
|
|
+ "fallback_count": row["fallback_count"],
|
|
|
|
|
+ "tool_count": row["tool_count"],
|
|
|
|
|
+ "turn_wall_time_ms": row["turn_wall_time_ms"],
|
|
|
}
|
|
}
|
|
|
for row in turn_rows
|
|
for row in turn_rows
|
|
|
]
|
|
]
|
|
|
session = self._fetchone(
|
|
session = self._fetchone(
|
|
|
"""
|
|
"""
|
|
|
SELECT
|
|
SELECT
|
|
|
- COALESCE(SUM(prompt_tokens), 0) AS prompt_tokens,
|
|
|
|
|
- COALESCE(SUM(completion_tokens), 0) AS completion_tokens,
|
|
|
|
|
- COALESCE(SUM(total_tokens), 0) AS total_tokens,
|
|
|
|
|
- COALESCE(SUM(cached_tokens), 0) AS cached_tokens,
|
|
|
|
|
- COALESCE(SUM(elapsed_ms), 0) AS elapsed_ms
|
|
|
|
|
|
|
+ COALESCE(SUM(CASE WHEN call_kind != 'tool_execution'
|
|
|
|
|
+ THEN prompt_tokens ELSE 0 END), 0) AS prompt_tokens,
|
|
|
|
|
+ COALESCE(SUM(CASE WHEN call_kind != 'tool_execution'
|
|
|
|
|
+ THEN completion_tokens ELSE 0 END), 0) AS completion_tokens,
|
|
|
|
|
+ COALESCE(SUM(CASE WHEN call_kind != 'tool_execution'
|
|
|
|
|
+ THEN total_tokens ELSE 0 END), 0) AS total_tokens,
|
|
|
|
|
+ COALESCE(SUM(CASE WHEN call_kind != 'tool_execution'
|
|
|
|
|
+ THEN cached_tokens ELSE 0 END), 0) AS cached_tokens,
|
|
|
|
|
+ COALESCE(SUM(elapsed_ms), 0) AS elapsed_ms,
|
|
|
|
|
+ COUNT(*) AS call_count,
|
|
|
|
|
+ COALESCE(SUM(CASE WHEN call_kind = 'argument_fallback'
|
|
|
|
|
+ THEN 1 ELSE 0 END), 0) AS fallback_count,
|
|
|
|
|
+ COALESCE(SUM(CASE WHEN call_kind = 'tool_execution'
|
|
|
|
|
+ THEN 1 ELSE 0 END), 0) AS tool_count
|
|
|
FROM usage_stats
|
|
FROM usage_stats
|
|
|
WHERE session_id = ?
|
|
WHERE session_id = ?
|
|
|
""",
|
|
""",
|
|
|
(session_id,),
|
|
(session_id,),
|
|
|
)
|
|
)
|
|
|
|
|
+ wall_time = self._fetchone(
|
|
|
|
|
+ """
|
|
|
|
|
+ SELECT SUM(wall_time_ms) AS turn_wall_time_ms
|
|
|
|
|
+ FROM turns
|
|
|
|
|
+ WHERE session_id = ?
|
|
|
|
|
+ """,
|
|
|
|
|
+ (session_id,),
|
|
|
|
|
+ )
|
|
|
return {
|
|
return {
|
|
|
"calls": calls,
|
|
"calls": calls,
|
|
|
"turns": turns,
|
|
"turns": turns,
|
|
@@ -413,10 +474,14 @@ class SQLiteSessionStore:
|
|
|
"total_tokens": session["total_tokens"],
|
|
"total_tokens": session["total_tokens"],
|
|
|
"cached_tokens": session["cached_tokens"],
|
|
"cached_tokens": session["cached_tokens"],
|
|
|
"elapsed_ms": session["elapsed_ms"],
|
|
"elapsed_ms": session["elapsed_ms"],
|
|
|
|
|
+ "call_count": session["call_count"],
|
|
|
|
|
+ "fallback_count": session["fallback_count"],
|
|
|
|
|
+ "tool_count": session["tool_count"],
|
|
|
|
|
+ "turn_wall_time_ms": wall_time["turn_wall_time_ms"],
|
|
|
},
|
|
},
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
- def _usage_row(self, row: sqlite3.Row) -> dict[str, Any]:
|
|
|
|
|
|
|
+ def _usage_row(self, row: sqlite3.Row, *, mode: str) -> dict[str, Any]:
|
|
|
return {
|
|
return {
|
|
|
"id": row["id"],
|
|
"id": row["id"],
|
|
|
"turn_index": row["turn_index"],
|
|
"turn_index": row["turn_index"],
|
|
@@ -427,6 +492,14 @@ class SQLiteSessionStore:
|
|
|
"cached_tokens": row["cached_tokens"],
|
|
"cached_tokens": row["cached_tokens"],
|
|
|
"ttft_ms": row["ttft_ms"],
|
|
"ttft_ms": row["ttft_ms"],
|
|
|
"elapsed_ms": row["elapsed_ms"],
|
|
"elapsed_ms": row["elapsed_ms"],
|
|
|
|
|
+ "mode": mode,
|
|
|
|
|
+ "agent": row["agent"],
|
|
|
|
|
+ "call_kind": row["call_kind"],
|
|
|
|
|
+ "event_id": row["event_id"],
|
|
|
|
|
+ "event_name": row["event_name"],
|
|
|
|
|
+ "used_fallback": bool(row["used_fallback"]),
|
|
|
|
|
+ "tool_latency_ms": row["tool_latency_ms"],
|
|
|
|
|
+ "metric_key": row["metric_key"],
|
|
|
"created_at": row["created_at"],
|
|
"created_at": row["created_at"],
|
|
|
}
|
|
}
|
|
|
|
|
|
|
@@ -476,6 +549,7 @@ class SQLiteSessionStore:
|
|
|
user_message TEXT NOT NULL,
|
|
user_message TEXT NOT NULL,
|
|
|
started_at TEXT NOT NULL,
|
|
started_at TEXT NOT NULL,
|
|
|
completed_at TEXT,
|
|
completed_at TEXT,
|
|
|
|
|
+ wall_time_ms INTEGER,
|
|
|
PRIMARY KEY (session_id, turn_index)
|
|
PRIMARY KEY (session_id, turn_index)
|
|
|
);
|
|
);
|
|
|
|
|
|
|
@@ -512,6 +586,13 @@ class SQLiteSessionStore:
|
|
|
cached_tokens INTEGER NOT NULL,
|
|
cached_tokens INTEGER NOT NULL,
|
|
|
ttft_ms INTEGER,
|
|
ttft_ms INTEGER,
|
|
|
elapsed_ms INTEGER NOT NULL,
|
|
elapsed_ms INTEGER NOT NULL,
|
|
|
|
|
+ agent TEXT NOT NULL DEFAULT 'chat_agent',
|
|
|
|
|
+ call_kind TEXT NOT NULL DEFAULT 'chat_completion',
|
|
|
|
|
+ event_id TEXT,
|
|
|
|
|
+ event_name TEXT,
|
|
|
|
|
+ used_fallback INTEGER NOT NULL DEFAULT 0,
|
|
|
|
|
+ tool_latency_ms INTEGER,
|
|
|
|
|
+ metric_key TEXT,
|
|
|
created_at TEXT NOT NULL
|
|
created_at TEXT NOT NULL
|
|
|
);
|
|
);
|
|
|
"""
|
|
"""
|
|
@@ -532,6 +613,35 @@ class SQLiteSessionStore:
|
|
|
or "tool_calls_json" not in self._message_columns()
|
|
or "tool_calls_json" not in self._message_columns()
|
|
|
):
|
|
):
|
|
|
raise
|
|
raise
|
|
|
|
|
+ self._add_column_if_missing(
|
|
|
|
|
+ "turns",
|
|
|
|
|
+ "wall_time_ms",
|
|
|
|
|
+ "INTEGER",
|
|
|
|
|
+ )
|
|
|
|
|
+ for column_name, declaration in (
|
|
|
|
|
+ ("agent", "TEXT NOT NULL DEFAULT 'chat_agent'"),
|
|
|
|
|
+ ("call_kind", "TEXT NOT NULL DEFAULT 'chat_completion'"),
|
|
|
|
|
+ ("event_id", "TEXT"),
|
|
|
|
|
+ ("event_name", "TEXT"),
|
|
|
|
|
+ ("used_fallback", "INTEGER NOT NULL DEFAULT 0"),
|
|
|
|
|
+ ("tool_latency_ms", "INTEGER"),
|
|
|
|
|
+ ("metric_key", "TEXT"),
|
|
|
|
|
+ ):
|
|
|
|
|
+ self._add_column_if_missing(
|
|
|
|
|
+ "usage_stats",
|
|
|
|
|
+ column_name,
|
|
|
|
|
+ declaration,
|
|
|
|
|
+ )
|
|
|
|
|
+ self._connection.execute(
|
|
|
|
|
+ "DROP INDEX IF EXISTS usage_stats_metric_key_uq"
|
|
|
|
|
+ )
|
|
|
|
|
+ self._connection.execute(
|
|
|
|
|
+ """
|
|
|
|
|
+ CREATE UNIQUE INDEX IF NOT EXISTS usage_stats_metric_key_uq
|
|
|
|
|
+ ON usage_stats(session_id, metric_key)
|
|
|
|
|
+ WHERE metric_key IS NOT NULL
|
|
|
|
|
+ """
|
|
|
|
|
+ )
|
|
|
self._connection.commit()
|
|
self._connection.commit()
|
|
|
except BaseException:
|
|
except BaseException:
|
|
|
self._connection.rollback()
|
|
self._connection.rollback()
|
|
@@ -544,6 +654,26 @@ class SQLiteSessionStore:
|
|
|
for row in self._connection.execute("PRAGMA table_info(messages)")
|
|
for row in self._connection.execute("PRAGMA table_info(messages)")
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+ def _table_columns(self, table_name: str) -> set[str]:
|
|
|
|
|
+ assert self._connection is not None
|
|
|
|
|
+ return {
|
|
|
|
|
+ row[1]
|
|
|
|
|
+ for row in self._connection.execute(f"PRAGMA table_info({table_name})")
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ def _add_column_if_missing(
|
|
|
|
|
+ self,
|
|
|
|
|
+ table_name: str,
|
|
|
|
|
+ column_name: str,
|
|
|
|
|
+ declaration: str,
|
|
|
|
|
+ ) -> None:
|
|
|
|
|
+ assert self._connection is not None
|
|
|
|
|
+ if column_name in self._table_columns(table_name):
|
|
|
|
|
+ return
|
|
|
|
|
+ self._connection.execute(
|
|
|
|
|
+ f"ALTER TABLE {table_name} ADD COLUMN {column_name} {declaration}"
|
|
|
|
|
+ )
|
|
|
|
|
+
|
|
|
def _touch_session_locked(self, session_id: str, now: str) -> None:
|
|
def _touch_session_locked(self, session_id: str, now: str) -> None:
|
|
|
self._connect().execute(
|
|
self._connect().execute(
|
|
|
"UPDATE sessions SET updated_at = ? WHERE id = ?",
|
|
"UPDATE sessions SET updated_at = ? WHERE id = ?",
|