mirror of
https://github.com/turnstonelabs/turnstone.git
synced 2026-08-12 23:12:23 -06:00
98e96ab5f3
* feat(models): add per-alias concurrency admission Add registry-backed FIFO admission limits with queue-aware deadlines and full-stream leases. Expose max_concurrency through storage, admin configuration, OpenAPI, documentation, and diagrams, with role and live backend count coverage. * fix(api): omit null concurrency schema default Keep max_concurrency optional for presence-keyed updates without advertising a null default for its non-null integer OpenAPI shape.
526 lines
14 KiB
Plaintext
526 lines
14 KiB
Plaintext
@startuml
|
|
!theme plain
|
|
title Turnstone — Core Engine Classes
|
|
|
|
skinparam classAttributeIconSize 0
|
|
|
|
' SessionUI Protocol
|
|
interface "SessionUI" as SessionUI <<Protocol>> {
|
|
+ on_thinking_start()
|
|
+ on_thinking_stop()
|
|
+ on_reasoning_token(text: str)
|
|
+ on_content_token(text: str)
|
|
+ on_stream_end()
|
|
+ approve_tools(items: list) → (bool, str|None)
|
|
+ on_tool_result(call_id: str, name: str, output: str, *, is_error: bool = False)
|
|
+ on_tool_output_chunk(call_id: str, chunk: str)
|
|
+ on_status(usage: dict, ctx_window: int, effort: str)
|
|
+ on_info(message: str)
|
|
+ on_error(message: str)
|
|
+ on_state_change(state: str)
|
|
+ on_rename(name: str)
|
|
}
|
|
|
|
' Implementations
|
|
class "TerminalUI" as TerminalUI {
|
|
Writes to stdout with ANSI colors
|
|
Prompts for approval via input()
|
|
--
|
|
cli.py
|
|
}
|
|
|
|
class "WorkstreamTerminalUI" as WsTermUI {
|
|
- _output_buffer: list[tuple]
|
|
- ws_id: str
|
|
- manager: SessionManager
|
|
+ flush_buffer()
|
|
--
|
|
Buffers output when workstream
|
|
is not foregrounded
|
|
}
|
|
|
|
class "WebUI" as WebUI {
|
|
- _listeners: list[Queue]
|
|
- _approval_cycles: dict[str, ApprovalCycle]
|
|
- _ws_prompt_tokens: int
|
|
- _ws_tool_calls: dict
|
|
+ resolve_approval(approved, feedback, cycle_id?, call_id?)
|
|
--
|
|
Enqueues JSON events for SSE.
|
|
Concurrent approval cycles each own
|
|
a threading.Event and result slot.
|
|
SSE handlers bridge Queue to
|
|
async via run_in_executor().
|
|
--
|
|
server.py
|
|
}
|
|
|
|
class "NullUI" as NullUI {
|
|
approve_tools() → (True, None)
|
|
All other methods: no-op
|
|
--
|
|
eval.py
|
|
}
|
|
|
|
' LLMProvider Protocol
|
|
interface "LLMProvider" as LLMProvider <<Protocol>> {
|
|
+ provider_name: str {property}
|
|
+ get_capabilities(model) → ModelCapabilities
|
|
+ create_streaming(client, model, messages, ..., cancel_ref, replay_reasoning_to_model) → Iterator[StreamChunk]
|
|
+ convert_tools(tools) → list[dict]
|
|
+ extract_reasoning_text(provider_blocks) → str
|
|
+ retryable_error_names: frozenset[str] {property}
|
|
--
|
|
core/providers/_protocol.py
|
|
}
|
|
|
|
class "OpenAIProvider" as OpenAIProv {
|
|
Model capability lookup table
|
|
(GPT-5.x, O-series, search)
|
|
Passthrough: messages already
|
|
in OpenAI format.
|
|
Search models: web_search_options
|
|
+ url_citation annotations.
|
|
Extended cache: 24h retention
|
|
for GPT-5.x (free).
|
|
--
|
|
core/providers/_openai.py
|
|
}
|
|
|
|
class "AnthropicProvider" as AnthropicProv {
|
|
Converts OpenAI messages to
|
|
Anthropic content blocks.
|
|
Adaptive + manual thinking.
|
|
Native web search via
|
|
web_search_20250305 server tool.
|
|
Auto prompt caching via
|
|
cache_control: ephemeral.
|
|
Lazy anthropic SDK import.
|
|
--
|
|
core/providers/_anthropic.py
|
|
}
|
|
|
|
class "GoogleProvider" as GoogleProv {
|
|
+ provider_name: str
|
|
+ get_capabilities(model) -> ModelCapabilities
|
|
--
|
|
Extends OpenAIChatCompletionsProvider
|
|
for Gemini /v1beta/openai/ endpoint.
|
|
Single default ModelCapabilities
|
|
(2M context, 65K output).
|
|
--
|
|
core/providers/_google.py
|
|
}
|
|
|
|
' ModelCapabilities
|
|
class "ModelCapabilities" as ModelCaps <<frozen>> {
|
|
+ context_window: int
|
|
+ max_output_tokens: int
|
|
+ supports_temperature: bool
|
|
+ token_param: str
|
|
+ thinking_mode: str
|
|
+ supports_effort: bool
|
|
+ supports_web_search: bool
|
|
+ supports_tool_search: bool
|
|
+ supports_vision: bool
|
|
+ supports_reasoning_replay: bool
|
|
}
|
|
|
|
class "ModelLane" as ModelLane <<frozen>> {
|
|
+ provider: LLMProvider
|
|
+ client: Any
|
|
+ model: str
|
|
+ alias: str
|
|
+ capabilities: ModelCapabilities
|
|
+ extra_params: dict | None
|
|
+ registry: ModelRegistry | None
|
|
+ admission: ModelAdmission | None
|
|
+ backend_auth_config: ModelConfig | None
|
|
+ backend_auth_resolver: Callable | None
|
|
}
|
|
|
|
class "ResolvedModelBinding" as ResolvedBinding <<frozen>> {
|
|
+ lane: ModelLane
|
|
+ config: ModelConfig | None
|
|
+ registry_generation: int
|
|
}
|
|
|
|
class "ModelTurnResult" as ModelTurnResult <<frozen>> {
|
|
+ turn: Turn
|
|
+ tool_calls: list[dict]
|
|
+ finish_reason: str
|
|
+ usage: UsageInfo | None
|
|
+ wire_msgs: list[dict] | None
|
|
+ producer: str
|
|
+ serving_model: str
|
|
}
|
|
|
|
class "model_turn()" as ModelTurnFn {
|
|
Turn IR → lower → provider stream
|
|
→ drain → canonical assistant Turn
|
|
--
|
|
core/model_turn.py
|
|
}
|
|
|
|
class "Backend auth resolver" as BackendAuth {
|
|
+ resolve_model_backend_auth_token(...)
|
|
--
|
|
Resolves static / Entra OBO /
|
|
Entra app / RFC 8693 per call.
|
|
Dynamic failure can fail closed.
|
|
--
|
|
core/model_backend_auth.py
|
|
}
|
|
|
|
' ChatSession
|
|
class "ChatSession" as ChatSession {
|
|
- _model_binding: ResolvedModelBinding
|
|
- _model_binding_lock: Lock
|
|
- ui: SessionUI
|
|
- messages: list[Turn]
|
|
- _msg_tokens: list[int]
|
|
- _ws_id: str
|
|
- _mcp_client: MCPClientManager | None
|
|
- _tool_search: ToolSearchManager | None
|
|
- _registry: ModelRegistry | None
|
|
- _generation: int
|
|
- _cancel_event: Event
|
|
- _durability_next_ticket: int
|
|
+ model_alias: str | None {property}
|
|
- _tools: list[dict]
|
|
- _task_tools: list[dict]
|
|
- _read_files: set[str]
|
|
- system_messages: list[dict]
|
|
--
|
|
+ send(user_input: str, ..., acting_user_id: str | None)
|
|
+ cancel()
|
|
+ compact_now() → bool
|
|
+ fork_from_storage(source_ws_id, principal_id, ...)
|
|
+ handle_command(command: str)
|
|
+ resume(ws_id: str)
|
|
- _save_config()
|
|
- _stream_response(my_generation) → ModelTurnResult
|
|
- _model_turn_with_fallback(consumer, prepare_wire) → ModelTurnResult
|
|
- _model_turn_with_retry(lane, tracker, ...) → ModelTurnResult
|
|
- _execute_tools(tool_calls) → (results, feedback)
|
|
- _prepare_tool(tc) → item dict
|
|
- _prepare_mcp_tool(call_id, name, args) → item dict
|
|
- _exec_mcp_tool(item) → (call_id, output)
|
|
- _get_active_tools() → list[dict]
|
|
- _prepare_tool_search() → None
|
|
- _exec_tool_search(item) → (call_id, output)
|
|
- _on_mcp_tools_changed()
|
|
- _rebuild_tool_search()
|
|
+ close()
|
|
- _run_agent(messages, tools, ...) → str
|
|
- _compact_messages(auto: bool, my_generation: int)
|
|
- _commit_for_generation(generation, commit)
|
|
- _publish_for_generation(generation, publish)
|
|
- _full_messages() → list[Turn]
|
|
- _update_token_table(msg)
|
|
- _emit_state(state: str)
|
|
- _generate_title()
|
|
}
|
|
|
|
' HeadlessSession
|
|
class "HeadlessSession" as HeadlessSession {
|
|
+ tool_call_log: list[dict]
|
|
+ auto_approve: bool = True
|
|
+ send_headless(input, max_turns, ...)
|
|
- _override_system_prompt(content)
|
|
--
|
|
eval.py: drained single-shot turns,
|
|
records all tool calls
|
|
}
|
|
|
|
' SessionManager
|
|
interface "SessionKindAdapter" as KindAdapter <<Protocol>> {
|
|
+ kind: WorkstreamKind
|
|
+ build_ui(ws) → SessionUI
|
|
+ build_session(ws, ...) → ChatSession
|
|
+ cleanup_ui(ws)
|
|
}
|
|
|
|
interface "SessionEventEmitter" as EventEmitter <<Protocol>> {
|
|
+ emit_created(ws)
|
|
+ emit_rehydrated(ws)
|
|
+ emit_state(ws, state)
|
|
+ emit_closed(ws_id, reason, name)
|
|
}
|
|
|
|
class "SessionManager" as SessionMgr {
|
|
- _adapter: SessionKindAdapter
|
|
- _storage: StorageBackend
|
|
- _workstreams: dict[str, Workstream]
|
|
- _pending_creates: dict[str, Workstream]
|
|
- _retiring_ids: set[str]
|
|
- _state_writer: StateWriter | None
|
|
- _order: list[str]
|
|
- _active_id: str
|
|
--
|
|
+ create(user_id, name, ..., defer_emit_created) → Workstream
|
|
+ commit_create(ws) → bool
|
|
+ discard(ws, ...) → bool
|
|
+ open(ws_id) → Workstream | None
|
|
+ delete(ws_id) → bool
|
|
+ close(ws_id)
|
|
+ get(ws_id) → Workstream
|
|
+ get_active() → Workstream
|
|
+ switch(ws_id)
|
|
+ set_state(ws_id, state)
|
|
+ close_idle(max_age_seconds)
|
|
- _max_workstreams: int
|
|
}
|
|
|
|
' Workstream
|
|
class "Workstream" as Ws <<dataclass>> {
|
|
+ id: str
|
|
+ name: str
|
|
+ state: WorkstreamState
|
|
+ session: ChatSession | None
|
|
+ ui: SessionUI | None
|
|
+ worker_thread: Thread | None
|
|
+ error_message: str
|
|
+ last_active: float
|
|
+ kind: WorkstreamKind
|
|
+ user_id: str
|
|
+ parent_ws_id: str | None
|
|
+ project_id: str | None
|
|
- _fork_reservation_token: str
|
|
- _closed: bool
|
|
- _lock: Lock
|
|
}
|
|
|
|
' WorkstreamState
|
|
enum "WorkstreamState" as WsState {
|
|
IDLE
|
|
THINKING
|
|
RUNNING
|
|
ATTENTION
|
|
ERROR
|
|
}
|
|
|
|
' MCPClientManager
|
|
class "MCPClientManager" as MCPMgr {
|
|
- _sessions: dict[str, ClientSession]
|
|
- _per_server_tools: dict[str, list[dict]]
|
|
- _per_server_resources: dict[str, list[dict]]
|
|
- _per_server_prompts: dict[str, list[dict]]
|
|
- _tools: list[dict]
|
|
- _tool_map: dict[str, tuple]
|
|
- _resource_map: dict[str, tuple]
|
|
- _prompt_map: dict[str, tuple]
|
|
- _supports_list_changed: dict[str, bool]
|
|
- _listeners: list[Callable]
|
|
--
|
|
+ start()
|
|
+ get_tools() → list[dict]
|
|
+ get_resources() → list[dict]
|
|
+ get_prompts() → list[dict]
|
|
+ is_mcp_tool(name) → bool
|
|
+ call_tool_sync(name, args) → str
|
|
+ read_resource_sync(uri) → str
|
|
+ get_prompt_sync(name, args?) → list[dict]
|
|
+ refresh_sync(server?) → dict
|
|
+ add_listener(callback)
|
|
+ remove_listener(callback)
|
|
+ server_names: list[str] {property}
|
|
+ shutdown()
|
|
--
|
|
Background asyncio event loop
|
|
bridges async MCP SDK to
|
|
sync ChatSession dispatch.
|
|
Push + manual refresh.
|
|
Resources + prompts discovered
|
|
alongside tools at startup.
|
|
--
|
|
core/mcp_client.py
|
|
}
|
|
|
|
' ToolSearchManager
|
|
class "ToolSearchManager" as ToolSearchMgr {
|
|
- _always_on: list[dict]
|
|
- _deferred: list[dict]
|
|
- _expanded: dict[str, None]
|
|
- _index: BM25Index
|
|
--
|
|
+ get_visible_tools() → list[dict]
|
|
+ get_deferred_tools() → list[dict]
|
|
+ get_expanded_names() → list[str]
|
|
+ search(query, k) → list[dict]
|
|
+ expand_visible(names) → list[dict]
|
|
+ get_search_tool_definition() → dict
|
|
+ format_search_results(tools) → str
|
|
}
|
|
|
|
' ModelRegistry
|
|
class "ModelRegistry" as ModelReg {
|
|
- _models: dict[str, ModelConfig]
|
|
- _clients: dict[str, Any]
|
|
- _providers: dict[str, LLMProvider]
|
|
- _admissions: dict[str, ModelAdmission]
|
|
- _client_lock: Lock
|
|
+ default: str
|
|
+ fallback: list[str]
|
|
+ agent_model: str | None
|
|
--
|
|
+ resolve_binding(alias) → (client, model, config, provider, admission, generation)
|
|
+ get_client(alias) → Any
|
|
+ get_provider(alias) → LLMProvider
|
|
+ has_alias(alias) → bool
|
|
+ list_aliases() → list[str]
|
|
+ shutdown()
|
|
--
|
|
Thread-safe lazy client + provider
|
|
creation. Loaded by load_model_registry()
|
|
from DB + [models.*] config + CLI args.
|
|
--
|
|
core/model_registry.py
|
|
}
|
|
|
|
class "ModelAdmission" as ModelAdmission {
|
|
- alias: str
|
|
- _limit: int
|
|
- _in_flight: int
|
|
- _waiters: deque
|
|
+ acquire(cancel_ref) → AdmissionLease
|
|
+ set_limit(limit)
|
|
+ snapshot() → AdmissionSnapshot
|
|
--
|
|
Per-process FIFO generation gate.
|
|
Stable across alias hot reloads;
|
|
queue time is deadline credit.
|
|
--
|
|
core/admission.py
|
|
}
|
|
|
|
class "ModelConfig" as ModelCfg <<frozen>> {
|
|
+ alias: str
|
|
+ provider: str
|
|
+ base_url: str
|
|
+ model: str
|
|
+ context_window: int
|
|
+ temperature: float | None
|
|
+ max_tokens: int | None
|
|
+ reasoning_effort: str | None
|
|
+ max_concurrency: int
|
|
+ auth_mode: str
|
|
+ obo_audience: str
|
|
+ obo_scopes: str
|
|
}
|
|
|
|
' Circuit breaker state
|
|
enum "CircuitState" as CircuitState {
|
|
CLOSED
|
|
OPEN
|
|
HALF_OPEN
|
|
}
|
|
|
|
' BackendHealthMonitor
|
|
class "BackendHealthMonitor" as HealthMon {
|
|
- _state: CircuitState
|
|
- _consecutive_failures: int
|
|
- _probe_interval: float
|
|
- _probe_timeout: float
|
|
- _threshold: int
|
|
- _cooldown: float
|
|
--
|
|
+ start()
|
|
+ stop()
|
|
+ record_success()
|
|
+ record_failure()
|
|
+ should_allow_request() → bool
|
|
--
|
|
Daemon thread probes
|
|
client.models.list() on interval.
|
|
--
|
|
core/healthcheck.py
|
|
}
|
|
|
|
' RateLimiter
|
|
class "RateLimiter" as RateLimiter {
|
|
- _buckets: dict[str, TokenBucket]
|
|
- _rate: float
|
|
- _burst: int
|
|
--
|
|
+ check(client_ip: str, path: str) → (bool, float)
|
|
--
|
|
Per-IP token bucket.
|
|
Returns (allowed, retry_after).
|
|
--
|
|
core/ratelimit.py
|
|
}
|
|
|
|
' TokenBucket
|
|
class "TokenBucket" as TokenBucket {
|
|
- _tokens: float
|
|
- _rate: float
|
|
- _burst: int
|
|
- _last_refill: float
|
|
--
|
|
+ consume() → bool
|
|
+ retry_after: float {property}
|
|
--
|
|
Single-IP bucket with
|
|
refill rate and burst cap.
|
|
}
|
|
|
|
' Relationships
|
|
SessionUI <|.. TerminalUI
|
|
TerminalUI <|-- WsTermUI
|
|
SessionUI <|.. WebUI
|
|
SessionUI <|.. NullUI
|
|
|
|
LLMProvider <|.. OpenAIProv
|
|
LLMProvider <|.. AnthropicProv
|
|
OpenAIProv <|-- GoogleProv
|
|
|
|
ChatSession --> SessionUI : uses
|
|
ChatSession --> ResolvedBinding : owns coherent snapshot
|
|
ChatSession --> ModelTurnFn : every model-backed role
|
|
ChatSession --> MCPMgr : optional
|
|
ChatSession --o ToolSearchMgr : _tool_search
|
|
ChatSession --> ModelReg : optional
|
|
ChatSession <|-- HeadlessSession
|
|
|
|
SessionMgr --> "*" Ws : manages
|
|
SessionMgr --> KindAdapter : delegates construction
|
|
SessionMgr --> EventEmitter : lifecycle fan-out
|
|
Ws --> "1" ChatSession : wraps
|
|
Ws --> "1" SessionUI : wraps
|
|
Ws --> "1" WsState : has
|
|
|
|
KindAdapter ..> ChatSession : constructs
|
|
|
|
ModelReg --> "*" ModelCfg : holds
|
|
ModelReg --> "*" LLMProvider : caches
|
|
ModelReg --> "*" ModelAdmission : owns per alias
|
|
LLMProvider --> ModelCaps : returns
|
|
ModelReg --> ResolvedBinding : resolves atomically
|
|
ResolvedBinding --> ModelLane
|
|
ModelLane --> LLMProvider
|
|
ModelLane --> ModelCaps
|
|
ModelLane --> ModelCfg : auth/config snapshot
|
|
ModelLane --> ModelAdmission : admission lease
|
|
ModelTurnFn --> ModelLane
|
|
ModelTurnFn --> ModelTurnResult
|
|
ModelTurnFn ..> BackendAuth : per-call resolver
|
|
|
|
ChatSession --> HealthMon : checks circuit
|
|
HealthMon --> "1" CircuitState : has
|
|
RateLimiter --> "*" TokenBucket : per IP
|
|
|
|
note bottom of ChatSession
|
|
Central engine: multi-turn LLM loop
|
|
with tool dispatch, agent sub-sessions,
|
|
context compaction, and memory persistence.
|
|
Provider-agnostic — delegates all LLM
|
|
communication to LLMProvider adapters.
|
|
|
|
Every live/durable publication is fenced by
|
|
its generation. Model calls use immutable lanes;
|
|
provider-wire mutation stays at lowering.
|
|
end note
|
|
|
|
@enduml
|