diff --git a/.env.example b/.env.example index 0efc439d..3438ecdb 100644 --- a/.env.example +++ b/.env.example @@ -1,37 +1,37 @@ # flux environment variables — copy to .env and fill in -# Provider API keys (or store them in the OS keychain — see -# docs/guides/CREDENTIAL-SETUP-FLOW.md) -FLUX_API_KEY= -OPENAI_API_KEY= +# Provider credentials, one per registry entry in catalog/registry/providers.go +# (registry SortOrder). Or store them in the OS keychain — see +# docs/guides/CREDENTIAL-SETUP-FLOW.md. catalog/registry/docs_test.go keeps this +# block in sync with the registry. +AGNES_API_KEY= +AWS_SECRET_ACCESS_KEY= ANTHROPIC_API_KEY= -GEMINI_API_KEY= +AZURE_OPENAI_API_KEY= +CANOPYWAVE_API_KEY= +CLINE_API_KEY= +CONCENTRATE_API_KEY= DEEPSEEK_API_KEY= +GEMINI_API_KEY= GROQ_API_KEY= -KIMI_API_KEY= MOONSHOT_API_KEY= -ZAI_API_KEY= -ZAI_CODING_API_KEY= -XIAOMI_MIMO_PAYG_API_KEY= -XIAOMI_MIMO_TOKEN_PLAN_API_KEY= -MINIMAX_API_KEY= +LONGCAT_API_KEY= +MINIMAX_PAYG_API_KEY= MINIMAX_TOKEN_PLAN_API_KEY= -AZURE_OPENAI_API_KEY= -AWS_SECRET_ACCESS_KEY= -VERTEX_ACCESS_TOKEN= +OPENAI_API_KEY= +OPENCODEGO_API_KEY= OPENROUTER_API_KEY= -CONCENTRATE_API_KEY= -OPENGATEWAY_API_KEY= -STEPFUN_API_KEY= -AGNES_API_KEY= -LONGCAT_API_KEY= -FIREWORKS_API_KEY= -CANOPYWAVE_API_KEY= +OLLAMA_BASE_URL= POOLSIDE_API_KEY= -CLINE_API_KEY= -OPENCODEGO_API_KEY= +VERTEX_ACCESS_TOKEN= XAI_API_KEY= -OLLAMA_BASE_URL= +XIAOMI_MIMO_PAYG_API_KEY= +XIAOMI_MIMO_TOKEN_PLAN_API_KEY= +ZAI_CODING_API_KEY= +ZAI_API_KEY= +STEP_API_KEY= +OPENGATEWAY_API_KEY= +FIREWORKS_API_KEY= # Default model overrides OPENAI_MODEL=gpt-4o diff --git a/AGENTS.md b/AGENTS.md index a9080d98..f7e09cb7 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -8,7 +8,8 @@ When starting any new work (feature, fix, refactor, chore), always create a feat ## Design Principles -- **Model-agnostic** — single interface for 75+ LLM providers +- **Model-agnostic** — single interface for the 28 provider gateways in + `catalog/registry/providers.go` (see README "Supported Providers") - **Host-neutral engine** — Flux owns provider routing, transport, caching, retry/fallback, and normalized telemetry; hosts own product UX and semantics - **Streaming-first** — all responses are streamed; blocking is opt-in @@ -55,6 +56,12 @@ make ci # Full CI suite `StreamResult`, `ResponseFormat`, `ImageURLPart`, `InputAudioPart`) live in `llm` with no `engine` alias; widening the facade to cover them is a deliberate API change, not an incidental one. +- The facade already exposes a frozen set of engine-internal symbols + (`credentials.Store`/`MapStore`, `operationsgraph.Input`/`Export`, the + `provider/resilience` rate-limit config and `provider/cache.CacheConfig`), + listed in `docs/architecture/HOST-ENGINE-BOUNDARY.md`. Changing them breaks + hosts. `engine/host_surface_test.go` fails whenever that reachable set + changes; update its list and the doc together, deliberately. - `provider/core.Provider` is the lower-level provider contract; keep its method set stable and use it across feature packages - Streaming tests need careful goroutine management @@ -77,7 +84,7 @@ make ci # Full CI suite - **Provider interface**: `provider/core.Provider` with `Chat()`, `StreamChat()`, `Ping()`, `Name()` - **Core request types**: `provider/core.FluxMessage`, `FluxResponse`, `FluxTool`, `FluxUsage` - **Config struct**: `provider/core.FluxConfig` with `Provider`, `APIKey`, `BaseURL`, `Model`, `MaxRetries` -- **Provider implementations**: `provider/adapters/AnthropicClient`, `OpenAIClient`, `GeminiClient`, etc. +- **Provider implementations**: `provider/adapters/anthropic.go` (`AnthropicClient`), `openai.go` (`OpenAIClient`), `gemini.go` (`GeminiClient`), etc. - **Compatibility configs**: `provider/adapters.OpenAICompat`, `GrokCompat`, `OpenRouterCompat` - **Error type**: `FluxError` with `Provider`, `Op`, `StatusCode`, `RequestID`, `Message`, `Err` fields - **Stream types**: `StreamResult`, `SSEEvent`, `StreamEvent` — streaming is SSE-based @@ -140,14 +147,14 @@ make ci # Full CI suite | Azure provider | `provider/adapters/azure.go` | | Provider registry | `provider/adapters/provider_registry.go` | | Provider compatibility | `provider/adapters/compat.go` (`OpenAICompat`, `GrokCompat`, etc.) | -| SSE streaming | `provider/stream.go` (`parseSSEStream()`, `SSEEvent`) | +| SSE streaming | `provider/core/stream.go` (`parseSSEStream()`, `SSEEvent`) | | Retry logic | `provider/core/retry.go` (`RetryConfig`, `backoffDelay()`, `shouldRetry()`) | | Rate limiting | `provider/resilience/ratelimit.go`, `provider/resilience/adaptive_ratelimit.go` | | Caching | `provider/cache/cache.go`, `provider/cache/semantic_cache.go` | -| Fallback chains | `provider/resilience/fallback.go` | +| Fallback chains | `router/router.go` (fallback providers), `router/deployment_router.go` (fallback deployment stages) | | Auto-continuation | `provider/resilience/continuation.go` | -| Error types | `provider/errors.go` (`FluxError`, `IsRetriable()`, `IsAuthError()`) | -| Error constants | `errors/errors.go` (API error messages, prompt-too-long parsing) | +| Error types | `provider/core/errors.go` (`FluxError`, `IsRetriable()`, `IsAuthError()`) | +| Error constants | `types/errors.go` (API error messages, prompt-too-long parsing) | | Model catalog | `catalog/` (pricing, context windows, capabilities per provider) | | Credentials | `credentials/` (key storage, env detection, scrubbing) — `HasSecret` is silent on miss (boolean predicate); `LookupSecret` logs `Debug` on `ErrNotFound` and `Warn` on real backend errors | | Mock provider | `provider/testkit/mock.go` | diff --git a/README.md b/README.md index b3c8558f..839bc49d 100644 --- a/README.md +++ b/README.md @@ -61,6 +61,13 @@ Everything else is engine-internal: `provider`, `catalog`, `config`, contracts. Enforced by `rho/scripts/check-flux-engine-boundary.sh` and two Go AST tests in `rho/internal/testaudit/`. +Exception: a fixed set of engine-internal symbols is reachable through the +facade (for example `credentials.Store` behind `engine.Options.SecretStore` +and `operationsgraph.Input` behind `engine.OperationsGraphInput`). Those +symbols are frozen as part of the contract; see +[Frozen engine-internal types](docs/architecture/HOST-ENGINE-BOUNDARY.md#frozen-engine-internal-types). +`engine/host_surface_test.go` fails when that set changes. + - do not import `rho/internal/*` - do not import the removed legacy path `rho/shared/types` @@ -70,8 +77,9 @@ and two Go AST tests in `rho/internal/testaudit/`. go get github.com/GrayCodeAI/flux ``` -Requires Go 1.26+ and a configured provider credential. Minimal dependencies -(UUID, OpenTelemetry, SQLite, keyring). +Requires Go 1.26+ and a configured provider credential. Direct dependencies: +UUID, tiktoken tokenizer, OS keyring, OpenTelemetry, pure-Go SQLite, and gRPC +(linked only into `-tags grpc` builds of `internal/grpc`). ```go import ( @@ -171,9 +179,9 @@ Named `primary` / `weak` / `editor` model slots with fallback to primary, plus a `POST /rerank` endpoint (provider-backed with lexical fallback) and a `GET /ready` readiness probe alongside the existing health check. -### gRPC Skeleton +### gRPC Transport (opt-in, internal) -Dependency-free gRPC API skeleton behind the `grpc` build tag — wired when generated stubs are available. +`internal/grpc` holds an optional gRPC transport behind the `grpc` build tag. It serves `flux.v1.ChatService/Chat` with a registered `json` content subtype (no `.proto` files or generated stubs; clients call with `grpc.CallContentSubtype("json")`), backed by `EngineChatService` over `conversation.Engine`. The package is internal, so hosts cannot import it, and nothing in flux starts it. `google.golang.org/grpc` is a direct requirement in `go.mod`, so it appears in consumers' module graphs, but only `-tags grpc` builds link it. ## Documentation @@ -199,40 +207,40 @@ ANTHROPIC_API_KEY=sk-... go run ./examples/basic/ ## Supported Providers -28 provider gateways in `catalog/registry/providers.go` (rho `/config` uses the same list), listed in registry `SortOrder`: +28 provider gateways in `catalog/registry/providers.go` (rho `/config` uses the same list), listed in registry `SortOrder`. `catalog/registry/docs_test.go` fails when this table, the count, or `.env.example` drift from the registry. | Provider | ID | Env variable | |---|---|---| +| **Agnes** | `agnes` | `AGNES_API_KEY` | +| **Amazon Bedrock** | `bedrock` | `AWS_SECRET_ACCESS_KEY` (+ `AWS_ACCESS_KEY_ID`, `AWS_SESSION_TOKEN`) | | **Anthropic** | `anthropic` | `ANTHROPIC_API_KEY` | -| **OpenAI** | `openai` | `OPENAI_API_KEY` | -| **Google Gemini** | `gemini` | `GEMINI_API_KEY` | -| **DeepSeek** | `deepseek` | `DEEPSEEK_API_KEY` | -| **xAI (Grok)** | `grok` | `XAI_API_KEY` | -| **Kimi (Moonshot)** | `kimi` | `MOONSHOT_API_KEY` | -| **Z.AI — Coding Plan** | `zai_coding` | `ZAI_CODING_API_KEY` | -| **Z.AI — Pay-as-you-go** | `zai_payg` | `ZAI_API_KEY` | -| **Xiaomi (MiMo) Token Plan** | `xiaomi_mimo_token_plan` | `XIAOMI_MIMO_TOKEN_PLAN_API_KEY` (+ region `cn` / `sgp` / `ams`) | -| **Xiaomi (MiMo) Pay-as-you-go** | `xiaomi_mimo_payg` | `XIAOMI_MIMO_PAYG_API_KEY` | -| **MiniMax — Token Plan** | `minimax_token_plan` | `MINIMAX_TOKEN_PLAN_API_KEY` | -| **MiniMax — Pay-as-you-go** | `minimax_payg` | `MINIMAX_PAYG_API_KEY` | | **Azure OpenAI** | `azure` | `AZURE_OPENAI_API_KEY` (+ `AZURE_OPENAI_ENDPOINT`) | -| **Amazon Bedrock** | `bedrock` | `AWS_SECRET_ACCESS_KEY` (+ `AWS_ACCESS_KEY_ID`, `AWS_SESSION_TOKEN`) | -| **Vertex AI** | `vertex` | `VERTEX_ACCESS_TOKEN` (or `GOOGLE_OAUTH_ACCESS_TOKEN`) | -| **OpenRouter** | `openrouter` | `OPENROUTER_API_KEY` | | **CanopyWave** | `canopywave` | `CANOPYWAVE_API_KEY` | -| **Poolside** | `poolside` | `POOLSIDE_API_KEY` | -| **Groq** | `groq` | `GROQ_API_KEY` | | **ClinePass** | `clinepass` | `CLINE_API_KEY` | | **Concentrate** | `concentrate` | `CONCENTRATE_API_KEY` | -| **OpenGateway** | `opengateway` | `OPENGATEWAY_API_KEY` | -| **StepFun** | `stepfun` | `STEPFUN_API_KEY` | -| **Agnes** | `agnes` | `AGNES_API_KEY` | +| **DeepSeek** | `deepseek` | `DEEPSEEK_API_KEY` | +| **Google Gemini** | `gemini` | `GEMINI_API_KEY` | +| **Groq** | `groq` | `GROQ_API_KEY` | +| **Kimi (Moonshot)** | `kimi` | `MOONSHOT_API_KEY` | | **LongCat** | `longcat` | `LONGCAT_API_KEY` | -| **Fireworks AI** | `fireworks` | `FIREWORKS_API_KEY` | +| **MiniMax — Pay-as-you-go** | `minimax_payg` | `MINIMAX_PAYG_API_KEY` | +| **MiniMax — Token Plan** | `minimax_token_plan` | `MINIMAX_TOKEN_PLAN_API_KEY` | +| **OpenAI** | `openai` | `OPENAI_API_KEY` | | **OpenCode Go** | `opencodego` | `OPENCODEGO_API_KEY` | +| **OpenRouter** | `openrouter` | `OPENROUTER_API_KEY` | | **Ollama** | `ollama` | `OLLAMA_BASE_URL` (local; no API key) | +| **Poolside** | `poolside` | `POOLSIDE_API_KEY` | +| **Vertex AI** | `vertex` | `VERTEX_ACCESS_TOKEN` (or `GOOGLE_OAUTH_ACCESS_TOKEN`) | +| **xAI (Grok)** | `grok` | `XAI_API_KEY` | +| **Xiaomi (MiMo) Pay-as-you-go** | `xiaomi_mimo_payg` | `XIAOMI_MIMO_PAYG_API_KEY` | +| **Xiaomi (MiMo) Token Plan** | `xiaomi_mimo_token_plan` | `XIAOMI_MIMO_TOKEN_PLAN_API_KEY` (+ region `cn` / `sgp` / `ams`) | +| **Z.AI — Coding Plan** | `zai_coding` | `ZAI_CODING_API_KEY` (+ region `international` / `cn`) | +| **Z.AI — Pay-as-you-go** | `zai_payg` | `ZAI_API_KEY` (+ region `international` / `cn`) | +| **StepFun** | `stepfun` | `STEP_API_KEY` (+ region `global` / `cn`) | +| **OpenGateway** | `opengateway` | `OPENGATEWAY_API_KEY` | +| **Fireworks AI** | `fireworks` | `FIREWORKS_API_KEY` | -Runtime auto-detection uses a separate priority order for chat when no deployment is pinned; see `config` profiles. +Runtime auto-detection uses a separate priority order (`config.APIProviderDetectionOrder`) when no deployment is pinned. ## Usage @@ -289,38 +297,54 @@ config.SaveProviderConfig(cfg, "") // save changes ``` flux/ -├── engine/ # Stable host-facing facade and provider-neutral DTOs -├── provider/ # Provider runtime and feature packages +├── engine/ # Stable host-facing facade (hosts import engine, llm, graph, tools) +├── llm/ # Host-facing DTOs and the Provider port that engine re-exports +├── graph/ # Portable execution-graph vocabulary +├── tools/ # Tool-call and tool-result contracts +├── provider/ # Provider runtime composition root (FluxClient) │ ├── core/ # Provider-neutral wire, stream, retry, and transport primitives │ ├── adapters/ # Provider protocol adapters and construction registry -│ └── embeddings/ # Embedding clients, cache, and defaults -├── config/ # Provider configuration & routing -│ └── credential/ # Credential file management +│ ├── resilience/ # Rate limits, continuation, guardrails, and error policy +│ ├── cache/ # Response and semantic caches +│ ├── batch/ # Batch execution +│ ├── embeddings/ # Embedding clients, cache, and defaults +│ ├── media/ # Image and audio clients, structured prompts +│ ├── extraction/ # Structured extraction +│ ├── observability/ # Usage, cost, metrics, tracing, and recording +│ └── testkit/ # Mock provider for tests ├── catalog/ # Model catalog & tier system +│ ├── registry/ # Provider registry (single source of truth for providers) │ ├── discover/ # Model discovery -│ ├── legacy/ # Legacy model support -│ ├── live/ # Live model data -│ └── registry/ # Model registry -├── codeagent/ # Code agent retry & fallback strategies -├── conversation/ # Conversation engine with branching -├── credentials/ # Credential management -├── docs/ # Documentation & guides -├── examples/ # Runnable code examples -├── router/ # Provider routing strategies +│ ├── live/ # Live model listing per provider +│ ├── capabilities/ # Capability and deprecation data +│ └── concentrate/ opencodego/ opengateway/ xiaomi/ zai/ # Gateway-specific helpers +├── config/ # Provider configuration & routing +│ └── credential/ # Credential file management +├── credentials/ # Keyring/env credential stores and OIDC keyless auth +├── router/ # Routing strategies, deployment router, circuit breakers +│ └── controlplane/ # Versioned, signed peer manifests and replicas +├── runtime/ # Engine-internal provider/model/credential resolution +├── setup/ # Catalog-backed deployment wiring ├── operationsgraph/ # Privacy-safe route and generation telemetry projection -├── runtime/ # Runtime manifest & routing policies -├── storage/ # SQLite conversation DAG store -├── types/ # Branded types & API errors -├── errors/ # Error message constants +├── conversation/ # Conversation engine with branching +├── storage/ # SQLite conversation DAG store, virtual keys, budgets +├── codeagent/ # Code agent retry & fallback strategies +├── verify/ # Provider conformance harness +├── types/ # Shared message types & API errors ├── constants/ # API limits ├── utils/ # Error utilities +├── api/ # OpenAPI spec for internal/api ├── internal/ -│ ├── api/ # HTTP API handlers -│ ├── cache/ # Response cache warmer +│ ├── api/ # HTTP API server (library code; no flux binary starts it) +│ ├── cache/ # Cache backends and response cache warmer +│ ├── grpc/ # Optional gRPC transport (build tag grpc) │ ├── health/ # Provider health checker -│ ├── observability/ # OpenTelemetry spans & metrics -│ ├── sdk/ # Go, Python, TypeScript client SDKs -│ └── version/ # Version information +│ ├── httputil/ probehttp/ shrink/ # HTTP, probe, and tool-description helpers +│ ├── observability/ # OpenTelemetry spans, metrics, and audit sinks +│ └── sdk/ # Go, Python, TypeScript clients for the internal/api HTTP surface +├── docs/ # Documentation & guides +├── examples/ # Runnable code examples +├── scripts/ # CI guards and helper scripts └── assets/ # Logo and branding ``` diff --git a/catalog/registry/docs_test.go b/catalog/registry/docs_test.go new file mode 100644 index 00000000..0a73a4ab --- /dev/null +++ b/catalog/registry/docs_test.go @@ -0,0 +1,156 @@ +package registry_test + +import ( + "os" + "path/filepath" + "regexp" + "sort" + "strconv" + "strings" + "testing" + + "github.com/GrayCodeAI/flux/catalog/registry" +) + +// These tests keep the provider facts that humans and agents read first +// (README, AGENTS.md, docs, .env.example, package comments) derived from the +// registry instead of hand-maintained numbers that drift. + +// repoRoot is the module root relative to this package directory. +const repoRoot = "../.." + +// providerCountDocs are the files that may state how many providers Flux +// supports. Every count they state must equal len(registry.All()). +var providerCountDocs = []string{ + "README.md", + "AGENTS.md", + "docs/README.md", + "docs/ARCHITECTURE.md", + "docs/guides/CREDENTIAL-SETUP-FLOW.md", + "docs/guides/DYNAMIC-MODEL-DISCOVERY.md", + "runtime/runtime.go", +} + +// providerCountClaim matches prose such as "28 provider gateways", +// "16 registered providers" or "75+ LLM providers". +var providerCountClaim = regexp.MustCompile(`(?i)\b(\d+)(\+?)\s+(?:registered\s+|supported\s+|LLM\s+)?provider(?:s|\s+gateways)\b`) + +// readmeProviderRow matches one row of the README "Supported Providers" table: +// | **Display name** | `provider_id` | `CREDENTIAL_ENV` optional note | +var readmeProviderRow = regexp.MustCompile("^\\| \\*\\*[^|]+\\*\\* \\| `([a-z0-9_]+)` \\| `([A-Z0-9_]+)`([^|]*)\\|$") + +var envAssignment = regexp.MustCompile(`^([A-Z][A-Z0-9_]*)=`) + +func readRepoFile(t *testing.T, rel string) string { + t.Helper() + data, err := os.ReadFile(filepath.Join(repoRoot, filepath.FromSlash(rel))) + if err != nil { + t.Fatalf("read %s: %v", rel, err) + } + return string(data) +} + +func specsBySortOrder() []registry.ProviderSpec { + specs := registry.All() + sort.SliceStable(specs, func(i, j int) bool { return specs[i].SortOrder < specs[j].SortOrder }) + return specs +} + +// markdownSection returns the text from heading up to the next level-2 heading. +func markdownSection(t *testing.T, doc, heading string) string { + t.Helper() + start := strings.Index(doc, "\n"+heading+"\n") + if start < 0 { + t.Fatalf("heading %q not found", heading) + } + body := doc[start+len(heading)+2:] + if end := strings.Index(body, "\n## "); end >= 0 { + body = body[:end] + } + return body +} + +func TestDocumentedProviderCountsMatchRegistry(t *testing.T) { + t.Parallel() + want := len(registry.All()) + for _, rel := range providerCountDocs { + for _, m := range providerCountClaim.FindAllStringSubmatch(readRepoFile(t, rel), -1) { + n, err := strconv.Atoi(m[1]) + if err != nil || n != want || m[2] != "" { + t.Errorf("%s claims %q; catalog/registry defines exactly %d providers", rel, m[0], want) + } + } + } +} + +func TestREADMEProviderTableMatchesRegistry(t *testing.T) { + t.Parallel() + section := markdownSection(t, readRepoFile(t, "README.md"), "## Supported Providers") + var rows [][]string + for _, line := range strings.Split(section, "\n") { + if m := readmeProviderRow.FindStringSubmatch(strings.TrimSpace(line)); m != nil { + rows = append(rows, m) + } + } + specs := specsBySortOrder() + if len(rows) != len(specs) { + t.Fatalf("README Supported Providers table has %d rows; registry has %d providers", len(rows), len(specs)) + } + for i, spec := range specs { + id, env, note := rows[i][1], rows[i][2], rows[i][3] + if id != spec.ProviderID { + t.Errorf("README row %d is %q; registry SortOrder position %d is %q", i+1, id, i+1, spec.ProviderID) + continue + } + if env != spec.CredentialEnv { + t.Errorf("README row %q lists %s; registry CredentialEnv is %s", id, env, spec.CredentialEnv) + } + for _, fallback := range spec.CredentialEnvFallbacks { + if !strings.Contains(note, "`"+fallback+"`") { + t.Errorf("README row %q does not mention credential fallback %s", id, fallback) + } + } + for _, region := range spec.RegionOptions { + if !strings.Contains(note, "`"+region.Value+"`") { + t.Errorf("README row %q does not mention region %q", id, region.Value) + } + } + } +} + +func TestEnvExampleListsRegistryCredentials(t *testing.T) { + t.Parallel() + credentialEnvs := map[string]bool{} + var want []string + for _, spec := range specsBySortOrder() { + credentialEnvs[spec.CredentialEnv] = true + want = append(want, spec.CredentialEnv) + } + + // The first run of consecutive KEY= lines is the provider credential block. + var block []string + inBlock := false + for _, line := range strings.Split(readRepoFile(t, ".env.example"), "\n") { + m := envAssignment.FindStringSubmatch(line) + if m == nil { + if inBlock { + break + } + continue + } + inBlock = true + block = append(block, m[1]) + } + if strings.Join(block, " ") != strings.Join(want, " ") { + t.Errorf(".env.example credential block drifted from the registry (SortOrder)\n got: %v\nwant: %v", block, want) + } + + // No other API-key variable may appear: Flux reads none besides the + // registry credentials, so an extra one would be a key nothing uses. + for _, line := range strings.Split(readRepoFile(t, ".env.example"), "\n") { + if m := envAssignment.FindStringSubmatch(strings.TrimPrefix(strings.TrimSpace(line), "# ")); m != nil && + strings.HasSuffix(m[1], "_API_KEY") && !credentialEnvs[m[1]] { + t.Errorf(".env.example lists %s, which is not a registry CredentialEnv", m[1]) + } + } +} diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index 7d052c60..2bdb977d 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -5,8 +5,7 @@ **Universal LLM Provider Runtime** [![Go](https://img.shields.io/badge/Go-1.26+-00ADD8?logo=go)](https://go.dev/) -[![Port](https://img.shields.io/badge/Port-8080-orange)](https://www.iana.org/assignments/service-names-port-numbers/service-names-port-numbers.xhtml) -[![Protocol](https://img.shields.io/badge/Protocol-REST-blue)](https://swagger.io/specification/) +[![Library](https://img.shields.io/badge/type-Go%20library-blue)](https://pkg.go.dev/github.com/GrayCodeAI/flux) @@ -33,15 +32,15 @@ flux/ │ ├── adapters/ Provider wire-protocol adapters │ ├── batch/ cache/ Batch execution and response caches │ ├── embeddings/ media/ Embeddings and multimodal features -│ ├── resilience/ Retry, fallback, rate limits and health +│ ├── resilience/ Rate limits, continuation, guardrails, health, error policy │ └── observability/ Usage, metrics, tracing and recording ├── catalog/ Model catalog and capabilities ├── config/ + credentials/ Config + keyring/env credential resolution ├── router/ Deployment policy and instance-local circuit breakers │ └── controlplane/ Versioned, signed peer manifests and replicas -├── runtime/ Host-facing construction +├── runtime/ Engine-internal provider/model/credential resolution ├── conversation/ + storage/ Conversation graph (branching DAG) + SQLite store -└── internal/api|cache|health|observability HTTP server, cache, health, OTel +└── internal/api|cache|grpc|health|observability HTTP server, cache, gRPC, health, OTel ``` The current distributed-routing foundation and its limits are described in @@ -51,11 +50,20 @@ The current distributed-routing foundation and its limits are described in ## globe API +flux is a Go library: it ships no binary, no `cmd/`, and no `flux serve` +command. Hosts such as rho call the [`engine`](../engine/) facade in-process. + +`internal/api` contains an HTTP server (`api.NewServer(api.Config{...})`, then +`ListenAndServe(addr)`) whose contract is [`api/openapi.yaml`](../api/openapi.yaml). +The package is internal, so only code inside the flux module can construct it, +and nothing in flux does today. The OpenAPI `servers` entry +(`http://localhost:8080`) is an example address, not a default. + | | | |---|---| | **Contract** | [`api/openapi.yaml`](../api/openapi.yaml) | -| **Port** | `:8080` (default). Override: `flux serve ` | -| **Auth** | Bearer token or `X-API-Key` header. Set via `FLUX_API_KEY` | +| **Address** | Whatever the embedding code passes to `ListenAndServe(addr)`; there is no default port | +| **Auth** | `Authorization: Bearer ` or `X-API-Key: `, compared with `api.Config.APIKey` on every route except `/health` and `/ready`. With an empty key the server refuses to bind a non-loopback address. There is no environment variable for the key. |
radio Endpoint Summary @@ -78,37 +86,42 @@ The current distributed-routing foundation and its limits are described in | `POST` | `/rerank` | rerank | Provider rerank + lexical fallback | | `GET` | `/ready` | health | Readiness probe (vs `/health` liveness) | +The last three endpoints are served by `internal/api` but are not yet described +in `api/openapi.yaml`. +
--- ## search Provider Detection -Auto-detects active provider from env vars in priority order: - -| Priority | Env Var | Provider | -|:--------:|---------|----------| -| 1 | `ANTHROPIC_API_KEY` | circle Anthropic Claude | -| 2 | `OPENAI_API_KEY` | circle OpenAI | -| 3 | `GEMINI_API_KEY` | circle Google Gemini | -| 4 | `OPENROUTER_API_KEY` | shuffle OpenRouter | -| 5 | `CANOPYWAVE_API_KEY` | radio CanopyWave | -| 6 | `XAI_API_KEY` | zap Grok (xAI) | -| 7 | `ZAI_API_KEY` | bot ZAI | -| 8 | — | server Ollama (localhost socket) | - -*Top 8 shown; full 28 in `catalog/registry/providers.go` — see [`CREDENTIAL-SETUP-FLOW.md`](./guides/CREDENTIAL-SETUP-FLOW.md) and `config` ChatPreference order.* +`provider.DetectProvider()` walks `config.APIProviderDetectionOrder` +(`config/profiles.go`) and returns the first provider whose credentials are +present in the credential store (by default the OS secret store; hosts inject +their own through `engine.Options.SecretStore`), defaulting to `anthropic` when +none is found. Multi-field providers need every field: Azure +needs `AZURE_OPENAI_API_KEY` and `AZURE_OPENAI_ENDPOINT`, Bedrock needs +`AWS_ACCESS_KEY_ID` and `AWS_SECRET_ACCESS_KEY`, Vertex needs +`VERTEX_PROJECT_ID` and `VERTEX_ACCESS_TOKEN`, and Ollama is detected from +`OLLAMA_BASE_URL`. The detection order is separate from the registry +`SortOrder` used for display and from `ChatPreference`; every provider in +[`catalog/registry/providers.go`](../catalog/registry/providers.go) appears in +it. Hosts do not call `DetectProvider`; they select through `engine`. --- ## radio Streaming -All responses are streamed via **SSE**. Blocking responses wrap the stream internally. +Providers implement both a blocking `Chat` and an SSE-based `StreamChat`; +streamed provider events are normalized into `FluxStreamEvent`s. Hosts use +`engine.Stream` (pull-based, must be closed) or `engine.Generate`. Inside +flux, a `*provider.FluxClient` exposes the same pair: ```go -sr, err := provider.StreamChat(ctx, messages, opts) +sr, err := client.StreamChat(ctx, messages, opts) +if err != nil { ... } defer sr.Close() -for event := range sr.Events() { ... } +for event := range sr.Events { ... } ``` --- @@ -117,19 +130,22 @@ for event := range sr.Events() { ... } | Feature | Behavior | |---------|----------| -| **Retries** | HTTP 429, 500, 502, 503, 529 | +| **Retries** | HTTP 429, 500, 502, 503, 529 (`core.DefaultRetryConfig`) | | **Backoff** | Exponential + jitter | -| **Retry-After** | Respected on 429 responses | -| **Rate Limiting** | Per-provider token-bucket | +| **Retry-After** | Honored (seconds or HTTP date, capped at `MaxDelay`) on any retried response | +| **Rate Limiting** | Per-provider token bucket; optional adaptive limiter driven by rate-limit headers | --- ## database Caching -| Layer | Strategy | Key | -|-------|----------|-----| -| **Exact** | Hash match | provider + model + message hash | -| **Semantic** | Cosine similarity | Prompt embeddings (optional, configurable TTL) | +| Layer | Where | Strategy | Key | +|-------|-------|----------|-----| +| **Exact** | `provider/cache` (`CachedProvider`) | SHA-256 match, LRU + TTL | model, system prompt, temperature, messages | +| **Semantic** | `provider/embeddings` (`EmbeddingCachedProvider`) | Cosine similarity ≥ 0.95 by default, LRU + TTL | prompt embedding from a configured embedding model | + +Both layers are opt-in, skip requests above the temperature threshold, and +cache only blocking `Chat` responses; `StreamChat` passes through uncached. --- diff --git a/docs/README.md b/docs/README.md index 5cb7875f..765ed28d 100644 --- a/docs/README.md +++ b/docs/README.md @@ -11,7 +11,9 @@ Welcome to the Flux documentation. This directory contains detailed guides and r - **[Flux Enterprise](design/FLUX-ENTERPRISE.md)** — Enterprise surfaces - **[Provider Setup Guide](guides/CREDENTIAL-SETUP-FLOW.md)** — How to configure credentials and providers - **[Dynamic Model Discovery](guides/DYNAMIC-MODEL-DISCOVERY.md)** — Architecture and implementation details for live model discovery -- **[OpenAPI](../api/openapi.yaml)** — HTTP surface (`/v1/chat/completions`, `/rerank`, `/ready`, `/health`) +- **[Decentralized Flux routing](architecture/DECENTRALIZED-FLUX.md)** — Instance-local routing and signed peer manifests +- **[Feature-oriented architecture](architecture/FEATURE-MONOREPO.md)** — Package layout and layering rules +- **[OpenAPI](../api/openapi.yaml)** — Contract for the internal HTTP server in `internal/api` (health, prompt, nodes, aliases, analytics); `/v1/chat/completions`, `/rerank` and `/ready` are served but not yet in the spec. flux ships no binary that starts this server. ### Quick Links @@ -34,17 +36,23 @@ The [`examples/`](../examples/) directory contains runnable code samples: docs/ ├── README.md # This file ├── ARCHITECTURE.md # System architecture -├── architecture/HOST-ENGINE-BOUNDARY.md -├── design/FLUX-ENTERPRISE.md -├── api/openapi.yaml # POST /v1/chat/completions, POST /rerank, GET /ready -└── guides/ - ├── CREDENTIAL-SETUP-FLOW.md - ├── DYNAMIC-MODEL-DISCOVERY.md - ├── RETRY-FALLBACK.md # (planned) backoff, fallback chains, circuit breaker - ├── CACHING-AUDIT.md # (planned) cache backends, audit sinks - └── ROUTING-STRATEGIES.md # weighted, latency, cost-based +├── architecture/ +│ ├── HOST-ENGINE-BOUNDARY.md # Host contract, frozen engine-internal types +│ ├── DECENTRALIZED-FLUX.md # Instance-local routing, signed peer manifests +│ └── FEATURE-MONOREPO.md # Package layout and layering +├── design/ +│ └── FLUX-ENTERPRISE.md # Enterprise surfaces (design) +├── guides/ +│ ├── CREDENTIAL-SETUP-FLOW.md +│ └── DYNAMIC-MODEL-DISCOVERY.md +└── plans/ # Remediation plans and reviews ``` +The HTTP contract lives outside this directory, at +[`../api/openapi.yaml`](../api/openapi.yaml). Retry/fallback, caching and +routing strategies are described in [ARCHITECTURE.md](ARCHITECTURE.md); there +are no separate guides for them yet. + ## For Developers If you're contributing to Flux: @@ -72,5 +80,4 @@ API documentation is available at: ## Support - **Issues**: [GitHub Issues](https://github.com/GrayCodeAI/flux/issues) -- **Discussions**: [GitHub Discussions](https://github.com/GrayCodeAI/flux/discussions) - **Security**: See [SECURITY.md](../SECURITY.md) for vulnerability reporting diff --git a/docs/architecture/HOST-ENGINE-BOUNDARY.md b/docs/architecture/HOST-ENGINE-BOUNDARY.md index ba452775..ce98d24c 100644 --- a/docs/architecture/HOST-ENGINE-BOUNDARY.md +++ b/docs/architecture/HOST-ENGINE-BOUNDARY.md @@ -32,7 +32,7 @@ Provider model APIs The dependency is one-way: Flux must not import Rho. Rho's integration layer may import `flux/engine`; Rho command, conversation, and UI packages must not -assemble Flux's `catalog`, `client`, `config`, `credentials`, `router`, +assemble Flux's `catalog`, `provider`, `config`, `credentials`, `router`, `runtime`, or `setup` packages. ## Composition root @@ -77,6 +77,37 @@ stream events do not cross this boundary. `Model` keeps distinct `Owner`, `ProviderID`, `GatewayID`, `CanonicalID`, `Source`, and `LiveMetadata` fields so Rho does not reconstruct catalog meaning. +## Frozen engine-internal types + +Contract v2 is not closed over `engine`, `llm`, `graph`, and `tools`. The +engine-internal symbols below are reachable from the facade through aliases, +`Options` fields, and re-exported functions, so they are frozen as part of +contract v2. Changing their names, fields, method sets, or signatures breaks +hosts exactly like changing `engine` itself and needs a contract-version bump. + +| Frozen symbol | Reached through | +|---|---| +| `credentials.Store` | `Options.SecretStore`; `SetDefaultStore` / `DefaultStore` signatures | +| `credentials.MapStore` | alias `engine.MapStore` (test fixture) | +| `credentials.SetDefaultStore` | `engine.SetDefaultStore` (test fixture) | +| `credentials.DefaultStore` | `engine.DefaultStore` (test fixture) | +| `operationsgraph.Input` | alias `engine.OperationsGraphInput` | +| `operationsgraph.Export` | alias `engine.OperationsGraphExport` | +| `provider/resilience.AdaptiveRateLimitConfig` | `Options.RateLimitConfig` | +| `provider/resilience.HeaderExtractor` | `AdaptiveRateLimitConfig.HeaderExtractor` | +| `provider/resilience.RateLimitHeaders` | result of `HeaderExtractor` | +| `provider/cache.CacheConfig` | `Options.CacheConfig` | + +Rho uses `engine.MapStore`, `engine.DefaultStore`, `engine.SetDefaultStore`, +`engine.OperationsGraphInput`, and `engine.BuildOperationsGraph` today, and its +import-path boundary checks cannot see through the aliases. Flux therefore +guards the set itself: `engine/host_surface_test.go` walks the exported surface +of the four contract packages, follows every engine-internal symbol it reaches +through struct fields, signatures, and exported methods, and fails when the +reachable set differs from this table. To drop an entry, give `engine` its own +type and bump the contract version; to add one, update the test's list and +this table in the same change. + ## Credential-to-conversation flow ```text @@ -206,6 +237,7 @@ Rho with `GOWORK=off`. ## Compatibility policy Lower-level Flux packages remain public for non-Rho consumers and staged -migration, but they are not part of Rho's product boundary. Additive fields +migration, but they are not part of Rho's product boundary, apart from the +frozen symbols listed above. Additive fields and stream events are allowed within contract v2. Removing or changing stable DTO semantics requires a contract-version and semantic-version boundary. diff --git a/docs_paths_test.go b/docs_paths_test.go new file mode 100644 index 00000000..1b7baab1 --- /dev/null +++ b/docs_paths_test.go @@ -0,0 +1,117 @@ +package flux_test + +import ( + "os" + "path/filepath" + "regexp" + "strings" + "testing" + "unicode/utf8" +) + +// These tests keep the file paths that agents and contributors read first +// pointing at files that exist: every path-like token in AGENTS.md and every +// entry of the directory trees in README.md and docs/README.md. + +// agentsPathToken matches a backticked token that contains a slash, such as +// `provider/core/stream.go` or `provider/core.Provider`. +var agentsPathToken = regexp.MustCompile("`([A-Za-z0-9_.<>-]+(?:/[A-Za-z0-9_.<>-]*)+)`") + +// symbolRef matches the last path element of a package-qualified symbol such +// as core.Provider or adapters.OpenAICompat. +var symbolRef = regexp.MustCompile(`^([a-z0-9_]+)\.[A-Z][A-Za-z0-9_]*$`) + +// treeEntry matches one entry of a box-drawing directory tree, capturing the +// indentation and the names before an optional "# comment". +var treeEntry = regexp.MustCompile(`^((?:│ | )*)(?:├── |└── )([^#]+?)\s*(?:#.*)?$`) + +func TestAgentsMDPathsExist(t *testing.T) { + data, err := os.ReadFile("AGENTS.md") + if err != nil { + t.Fatal(err) + } + for _, m := range agentsPathToken.FindAllStringSubmatch(string(data), -1) { + token := m[1] + switch { + case strings.ContainsAny(token, "<>"), // placeholder such as provider/.go + strings.HasPrefix(token, "../"), // sibling checkout + strings.HasPrefix(token, "github.com/"), + strings.HasPrefix(token, "rho/"): + continue + } + path := strings.TrimSuffix(token, "/") + dir, last := filepath.Split(path) + if sm := symbolRef.FindStringSubmatch(last); sm != nil { + path = dir + sm[1] // package directory of a qualified symbol + } + if _, err := os.Stat(filepath.FromSlash(path)); err != nil { + t.Errorf("AGENTS.md cites `%s`, but %s does not exist", token, path) + } + } +} + +func TestDocTreesListExistingPaths(t *testing.T) { + for _, tc := range []struct{ doc, root string }{ + {"README.md", "flux/"}, + {"docs/README.md", "docs/"}, + } { + data, err := os.ReadFile(filepath.FromSlash(tc.doc)) + if err != nil { + t.Fatal(err) + } + checkTree(t, tc.doc, tc.root, string(data)) + } +} + +// checkTree finds the fenced tree whose first line is root and verifies that +// each entry exists. root "flux/" is the repository root. +func checkTree(t *testing.T, doc, root, text string) { + t.Helper() + lines := strings.Split(text, "\n") + start := -1 + for i, line := range lines { + if strings.TrimSpace(line) == root && i > 0 && strings.HasPrefix(lines[i-1], "```") { + start = i + 1 + break + } + } + if start < 0 { + t.Fatalf("%s: no directory tree rooted at %q", doc, root) + } + base := strings.TrimSuffix(root, "/") + if root == "flux/" { + base = "." + } + stack := []string{base} + entries := 0 + for _, line := range lines[start:] { + if strings.HasPrefix(line, "```") { + break + } + m := treeEntry.FindStringSubmatch(line) + if m == nil { + t.Errorf("%s: unparseable tree line %q", doc, line) + continue + } + depth := utf8.RuneCountInString(m[1]) / 4 + if depth+1 > len(stack) { + t.Errorf("%s: tree line %q is nested under nothing", doc, line) + continue + } + stack = stack[:depth+1] + names := strings.Fields(m[2]) + for _, name := range names { + entries++ + path := filepath.Join(append(append([]string{}, stack...), filepath.FromSlash(name))...) + if _, err := os.Stat(path); err != nil { + t.Errorf("%s lists %s, which does not exist", doc, filepath.ToSlash(path)) + } + } + if len(names) == 1 && strings.HasSuffix(names[0], "/") { + stack = append(stack, strings.TrimSuffix(names[0], "/")) + } + } + if entries == 0 { + t.Errorf("%s: tree rooted at %q has no entries", doc, root) + } +} diff --git a/engine/doc.go b/engine/doc.go index 35351fc1..3d18cbed 100644 --- a/engine/doc.go +++ b/engine/doc.go @@ -1,9 +1,15 @@ // Package engine is the stable, host-facing Flux API. // -// Hosts should prefer this package over assembling client, catalog, config, +// Hosts should prefer this package over assembling provider, catalog, config, // credentials, runtime, and setup packages directly. The lower-level packages // remain public for backward compatibility and advanced integrations. // +// A fixed set of engine-internal symbols is reachable from this package (for +// example Options.SecretStore is a credentials.Store and OperationsGraphInput +// aliases operationsgraph.Input). Those symbols are frozen as part of the +// contract; docs/architecture/HOST-ENGINE-BOUNDARY.md lists them and +// host_surface_test.go fails when the reachable set changes. +// // Engine is intentionally stateless with respect to product conversations: // the host owns conversation history, tools, permissions, and checkpoints; // Flux owns credential, catalog, selection, routing, and model transport. diff --git a/engine/host_surface_test.go b/engine/host_surface_test.go new file mode 100644 index 00000000..4a67f095 --- /dev/null +++ b/engine/host_surface_test.go @@ -0,0 +1,430 @@ +package engine_test + +import ( + "go/ast" + "go/build" + "go/parser" + "go/token" + "os" + "path/filepath" + "regexp" + "sort" + "strconv" + "strings" + "testing" +) + +// The host contract is the exported surface of engine, llm, graph and tools. +// Some of that surface is spelled with types from engine-internal packages +// (aliases such as engine.MapStore, Options fields such as SecretStore, and +// re-exported functions such as engine.SetDefaultStore). Those symbols are +// frozen: changing them breaks hosts exactly like changing engine itself. +// +// This test walks the exported surface of the four contract packages, follows +// every referenced engine-internal symbol transitively (struct fields, +// signatures, exported methods), and requires the reachable set to equal +// frozenInternalSymbols. A new leak, or a removed one, fails until the list +// and docs/architecture/HOST-ENGINE-BOUNDARY.md are updated deliberately. + +const fluxModule = "github.com/GrayCodeAI/flux" + +// fluxRoot is the module root relative to the engine package directory. +const fluxRoot = ".." + +var hostContractPackages = []string{ + fluxModule + "/engine", + fluxModule + "/graph", + fluxModule + "/llm", + fluxModule + "/tools", +} + +// frozenInternalSymbols is the complete set of engine-internal symbols +// reachable from the host contract. Keep it sorted and in sync with the +// "Frozen engine-internal types" section of HOST-ENGINE-BOUNDARY.md. +var frozenInternalSymbols = []string{ + fluxModule + "/credentials.DefaultStore", + fluxModule + "/credentials.MapStore", + fluxModule + "/credentials.SetDefaultStore", + fluxModule + "/credentials.Store", + fluxModule + "/operationsgraph.Export", + fluxModule + "/operationsgraph.Input", + fluxModule + "/provider/cache.CacheConfig", + fluxModule + "/provider/resilience.AdaptiveRateLimitConfig", + fluxModule + "/provider/resilience.HeaderExtractor", + fluxModule + "/provider/resilience.RateLimitHeaders", +} + +type surfaceDecl struct { + file *ast.File + typeSpec *ast.TypeSpec + funcDecl *ast.FuncDecl + valueSpec *ast.ValueSpec + index int // position of the name inside valueSpec +} + +type surfacePkg struct { + path string + name string + decls map[string]surfaceDecl + methods map[string][]surfaceDecl // receiver base type name -> methods + imports map[*ast.File]map[string]string +} + +type surfaceScanner struct { + t *testing.T + fset *token.FileSet + pkgs map[string]*surfacePkg + visited map[string]bool + internal map[string][]string // internal symbol -> contract paths that reach it +} + +func newSurfaceScanner(t *testing.T) *surfaceScanner { + return &surfaceScanner{ + t: t, + fset: token.NewFileSet(), + pkgs: map[string]*surfacePkg{}, + visited: map[string]bool{}, + internal: map[string][]string{}, + } +} + +func isHostContract(path string) bool { + for _, p := range hostContractPackages { + if p == path { + return true + } + } + return false +} + +func isFluxPackage(path string) bool { + return path == fluxModule || strings.HasPrefix(path, fluxModule+"/") +} + +// isStdlib mirrors the go command's rule: standard-library import paths have +// no dot in their first element. +func isStdlib(path string) bool { + first, _, _ := strings.Cut(path, "/") + return !strings.Contains(first, ".") +} + +var majorVersionElem = regexp.MustCompile(`^v[0-9]+$`) + +// assumedPackageName guesses the name of a non-Flux package imported without +// an explicit name, the same way goimports does. A wrong guess cannot hide a +// leak: an unresolved qualifier in a type expression fails the test. +func assumedPackageName(path string) string { + elems := strings.Split(path, "/") + name := elems[len(elems)-1] + if majorVersionElem.MatchString(name) && len(elems) > 1 { + name = elems[len(elems)-2] + } + if i := strings.Index(name, ".v"); i > 0 { + name = name[:i] + } + name = strings.TrimPrefix(name, "go-") + name = strings.TrimSuffix(name, "-go") + return strings.ReplaceAll(name, "-", "_") +} + +func (s *surfaceScanner) load(path string) *surfacePkg { + if pkg, ok := s.pkgs[path]; ok { + return pkg + } + dir := filepath.Join(fluxRoot, filepath.FromSlash(strings.TrimPrefix(path, fluxModule))) + bp, err := build.ImportDir(dir, 0) + if err != nil { + s.t.Fatalf("load %s: %v", path, err) + } + pkg := &surfacePkg{ + path: path, + name: bp.Name, + decls: map[string]surfaceDecl{}, + methods: map[string][]surfaceDecl{}, + imports: map[*ast.File]map[string]string{}, + } + s.pkgs[path] = pkg + for _, name := range bp.GoFiles { + file, err := parser.ParseFile(s.fset, filepath.Join(dir, name), nil, parser.SkipObjectResolution) + if err != nil { + s.t.Fatalf("parse %s/%s: %v", path, name, err) + } + for _, decl := range file.Decls { + switch d := decl.(type) { + case *ast.FuncDecl: + if d.Recv == nil { + pkg.decls[d.Name.Name] = surfaceDecl{file: file, funcDecl: d} + } else if len(d.Recv.List) == 1 { + recv := receiverTypeName(d.Recv.List[0].Type) + pkg.methods[recv] = append(pkg.methods[recv], surfaceDecl{file: file, funcDecl: d}) + } + case *ast.GenDecl: + for _, spec := range d.Specs { + switch sp := spec.(type) { + case *ast.TypeSpec: + pkg.decls[sp.Name.Name] = surfaceDecl{file: file, typeSpec: sp} + case *ast.ValueSpec: + for i, n := range sp.Names { + pkg.decls[n.Name] = surfaceDecl{file: file, valueSpec: sp, index: i} + } + } + } + } + } + } + return pkg +} + +func receiverTypeName(expr ast.Expr) string { + for { + switch e := expr.(type) { + case *ast.StarExpr: + expr = e.X + case *ast.IndexExpr: + expr = e.X + case *ast.IndexListExpr: + expr = e.X + case *ast.Ident: + return e.Name + default: + return "" + } + } +} + +func (s *surfaceScanner) importsOf(pkg *surfacePkg, file *ast.File) map[string]string { + if m, ok := pkg.imports[file]; ok { + return m + } + m := map[string]string{} + for _, imp := range file.Imports { + path, err := strconv.Unquote(imp.Path.Value) + if err != nil { + s.t.Fatalf("%s: bad import %s", pkg.path, imp.Path.Value) + } + var name string + switch { + case imp.Name != nil: + name = imp.Name.Name + case isFluxPackage(path): + name = s.load(path).name + default: + name = assumedPackageName(path) + } + if name == "." { + s.t.Errorf("%s: dot import of %s hides package qualifiers from this check", pkg.path, path) + continue + } + if name != "_" { + m[name] = path + } + } + pkg.imports[file] = m + return m +} + +// reference records that the contract surface reaches path.name and walks it. +func (s *surfaceScanner) reference(path, name, from string) { + switch { + case isHostContract(path): + return // scanned as a root + case isStdlib(path): + return + } + key := path + "." + name + s.internal[key] = append(s.internal[key], from) + if isFluxPackage(path) { + s.visit(s.load(path), name, key) + } +} + +// local handles an unqualified identifier declared in pkg. +func (s *surfaceScanner) local(pkg *surfacePkg, name, from string) { + if _, ok := pkg.decls[name]; !ok { + return // predeclared identifier or type parameter + } + if isHostContract(pkg.path) { + s.visit(pkg, name, from) + return + } + s.reference(pkg.path, name, from) +} + +func (s *surfaceScanner) visit(pkg *surfacePkg, name, from string) { + key := pkg.path + "." + name + if s.visited[key] { + return + } + s.visited[key] = true + d, ok := pkg.decls[name] + if !ok { + s.t.Errorf("%s references %s, which is not a top-level declaration", from, key) + return + } + switch { + case d.typeSpec != nil: + if d.typeSpec.TypeParams != nil { + s.walkType(pkg, d.file, d.typeSpec.TypeParams, key) + } + s.walkType(pkg, d.file, d.typeSpec.Type, key) + for _, m := range pkg.methods[name] { + if m.funcDecl.Name.IsExported() { + s.walkType(pkg, m.file, m.funcDecl.Type, key+"."+m.funcDecl.Name.Name) + } + } + case d.funcDecl != nil: + s.walkType(pkg, d.file, d.funcDecl.Type, key) + case d.valueSpec != nil: + switch { + case d.valueSpec.Type != nil: + s.walkType(pkg, d.file, d.valueSpec.Type, key) + case d.index < len(d.valueSpec.Values): + s.walkValue(pkg, d.file, d.valueSpec.Values[d.index], key) + } + // An untyped spec without values repeats the previous const spec, + // which is visited on its own. + } +} + +// walkType follows every named type in a type expression, skipping +// unexported struct fields, which are not part of the surface. +func (s *surfaceScanner) walkType(pkg *surfacePkg, file *ast.File, expr ast.Node, from string) { + ast.Inspect(expr, func(n ast.Node) bool { + switch n := n.(type) { + case *ast.StructType: + for _, f := range n.Fields.List { + if len(f.Names) > 0 && !anyExported(f.Names) { + continue + } + s.walkType(pkg, file, f.Type, from) + } + return false + case *ast.Field: + s.walkType(pkg, file, n.Type, from) + return false + case *ast.SelectorExpr: + id, ok := n.X.(*ast.Ident) + if !ok { + s.t.Errorf("%s: unexpected qualified type %T", from, n.X) + return false + } + path, ok := s.importsOf(pkg, file)[id.Name] + if !ok { + s.t.Errorf("%s: cannot resolve package qualifier %q; teach assumedPackageName about it", from, id.Name) + return false + } + s.reference(path, n.Sel.Name, from) + return false + case *ast.Ident: + s.local(pkg, n.Name, from) + return false + } + return true + }) +} + +// walkValue follows the declared type of an initializer for a var or const +// declared without an explicit type. +func (s *surfaceScanner) walkValue(pkg *surfacePkg, file *ast.File, expr ast.Expr, from string) { + switch e := expr.(type) { + case *ast.BasicLit: + case *ast.Ident: + s.local(pkg, e.Name, from) + case *ast.SelectorExpr: + if id, ok := e.X.(*ast.Ident); ok { + if path, ok := s.importsOf(pkg, file)[id.Name]; ok { + s.reference(path, e.Sel.Name, from) + return + } + } + s.t.Errorf("%s: cannot type initializer %s; declare the value with an explicit type", from, exprString(s.fset, e)) + case *ast.CallExpr: + s.walkValue(pkg, file, e.Fun, from) + case *ast.CompositeLit: + s.walkType(pkg, file, e.Type, from) + case *ast.FuncLit: + s.walkType(pkg, file, e.Type, from) + case *ast.ParenExpr: + s.walkValue(pkg, file, e.X, from) + case *ast.UnaryExpr: + s.walkValue(pkg, file, e.X, from) + case *ast.BinaryExpr: + s.walkValue(pkg, file, e.X, from) + s.walkValue(pkg, file, e.Y, from) + default: + s.t.Errorf("%s: unsupported initializer %s; declare the value with an explicit type", from, exprString(s.fset, e)) + } +} + +func exprString(fset *token.FileSet, e ast.Expr) string { + return fset.Position(e.Pos()).String() +} + +func anyExported(names []*ast.Ident) bool { + for _, n := range names { + if n.IsExported() { + return true + } + } + return false +} + +func (s *surfaceScanner) scanContract() { + for _, path := range hostContractPackages { + pkg := s.load(path) + names := make([]string, 0, len(pkg.decls)) + for name := range pkg.decls { + if ast.IsExported(name) { + names = append(names, name) + } + } + sort.Strings(names) + for _, name := range names { + s.visit(pkg, name, path+"."+name) + } + } +} + +func TestHostContractExposesOnlyFrozenInternalSymbols(t *testing.T) { + s := newSurfaceScanner(t) + s.scanContract() + + frozen := map[string]bool{} + for _, sym := range frozenInternalSymbols { + frozen[sym] = true + } + var reached []string + for sym := range s.internal { + reached = append(reached, sym) + } + sort.Strings(reached) + for _, sym := range reached { + if !frozen[sym] { + from := s.internal[sym] + sort.Strings(from) + t.Errorf("engine-internal symbol %s is reachable from the host contract via %s: stop exposing it, or freeze it in frozenInternalSymbols and HOST-ENGINE-BOUNDARY.md", sym, strings.Join(from, ", ")) + } + } + for _, sym := range frozenInternalSymbols { + if _, ok := s.internal[sym]; !ok { + t.Errorf("%s is frozen but no longer reachable from the host contract: remove it from frozenInternalSymbols and HOST-ENGINE-BOUNDARY.md", sym) + } + } + if !sort.StringsAreSorted(frozenInternalSymbols) { + t.Error("keep frozenInternalSymbols sorted") + } +} + +func TestFrozenInternalSymbolsAreDocumented(t *testing.T) { + data, err := os.ReadFile(filepath.Join(fluxRoot, "docs", "architecture", "HOST-ENGINE-BOUNDARY.md")) + if err != nil { + t.Fatal(err) + } + doc := string(data) + for _, sym := range frozenInternalSymbols { + short := "`" + strings.TrimPrefix(sym, fluxModule+"/") + "`" + if !strings.Contains(doc, short) { + t.Errorf("HOST-ENGINE-BOUNDARY.md does not list frozen symbol %s", short) + } + } +} diff --git a/internal/grpc/README.md b/internal/grpc/README.md index 86d3398a..3ff97957 100644 --- a/internal/grpc/README.md +++ b/internal/grpc/README.md @@ -5,9 +5,14 @@ Go request/response structs with the registered `json` gRPC content subtype, which avoids generated protobuf code while retaining gRPC framing, interceptors, deadlines, status propagation, and HTTP/2 transport. -- `grpc.go` defines the transport-independent `ChatService` contract. +- `grpc.go` defines the transport-independent `ChatService` contract and + `EngineChatService`, which serves a unary Chat as one + `conversation.Engine` prompt. - `server_grpc.go` registers and serves `flux.v1.ChatService/Chat`. - Clients must select `grpc.CallContentSubtype("json")`. +- There are no `.proto` files or generated stubs. +- The package is internal: hosts cannot import it, and nothing in Flux starts + the server. ## Running @@ -17,4 +22,6 @@ go build -tags grpc ./... Callers provide a `ChatService` implementation to `Serve` or `NewServer`. The untagged build retains only the service contract, so consumers that do not -need a network server do not link the gRPC runtime. +need a network server do not link the gRPC runtime. `google.golang.org/grpc` +is still a direct requirement in `go.mod` (the tagged file needs it), so it +appears in consumers' module graphs and `go.sum`. diff --git a/internal/grpc/grpc.go b/internal/grpc/grpc.go index 562a90f5..06cf0137 100644 --- a/internal/grpc/grpc.go +++ b/internal/grpc/grpc.go @@ -1,11 +1,17 @@ -// Package grpc holds a dependency-free skeleton for an flux gRPC API. +// Package grpc is Flux's optional gRPC transport for the conversation engine. // -// flux does not currently import google.golang.org/grpc, and per repo policy -// that dependency is not added speculatively. This file therefore defines only -// the service contract and a no-op default implementation so the rest of the -// codebase can reference the gRPC surface today. The real server wiring lives -// in server_grpc.go behind the "grpc" build tag. See README.md for the design -// note and codegen steps. +// This untagged file holds the transport-independent ChatService contract, +// its request/response structs, a no-op default (NewChatService) and +// EngineChatService, which serves a unary Chat as one conversation.Engine +// prompt. It does not import google.golang.org/grpc. +// +// server_grpc.go (build tag "grpc") imports google.golang.org/grpc, registers +// a "json" codec and serves flux.v1.ChatService/Chat. There are no .proto +// files or generated stubs: clients select grpc.CallContentSubtype("json"). +// Because of that tagged file, google.golang.org/grpc is a direct requirement +// in go.mod and appears in consumers' module graphs, although untagged builds +// do not link it. The package is internal, so hosts cannot import it, and +// nothing in Flux starts the server. See README.md. package grpc import ( @@ -34,21 +40,17 @@ type ChatResponse struct { } // ChatService is the flux gRPC service contract: a single unary Chat RPC. -// A concrete implementation will adapt conversation.Engine; see README.md. +// EngineChatService is the conversation.Engine-backed implementation. type ChatService interface { Chat(ctx context.Context, req *ChatRequest) (*ChatResponse, error) } -// noopChatService is the default ChatService. It returns ErrUnimplemented so -// callers get a clear signal that the gRPC backend has not been wired up. +// noopChatService is the placeholder ChatService. It returns ErrUnimplemented +// so callers get a clear signal that no backend was supplied. type noopChatService struct{} -// ErrUnimplemented is returned by the default ChatService until a real -// gRPC-backed implementation is provided. -// -// When google.golang.org/grpc and the generated protobuf stubs are added -// (see README.md), replace noopChatService with an engine-backed adapter -// and register it via server_grpc.go (build tag "grpc"). +// ErrUnimplemented is returned by the placeholder ChatService from +// NewChatService and by an EngineChatService built with a nil engine. var ErrUnimplemented = errUnimplemented{} type errUnimplemented struct{} @@ -59,9 +61,8 @@ func (noopChatService) Chat(_ context.Context, _ *ChatRequest) (*ChatResponse, e return nil, ErrUnimplemented } -// NewChatService returns the default (no-op) ChatService. It exists so callers -// have a stable constructor; once a real backend exists this will return the -// engine-backed implementation instead. +// NewChatService returns the placeholder (no-op) ChatService. Use +// NewEngineChatService for a working backend. func NewChatService() ChatService { return noopChatService{} } @@ -74,7 +75,8 @@ type EngineChatService struct { } // NewEngineChatService returns a ChatService backed by a conversation.Engine. -// It is the real backend referenced by the gRPC server (build tag "grpc"). +// Pass it to NewServer or Serve (build tag "grpc") to expose the engine over +// gRPC. func NewEngineChatService(engine *conversation.Engine) ChatService { return &EngineChatService{engine: engine} } diff --git a/internal/grpc/grpc_engine_test.go b/internal/grpc/grpc_engine_test.go index cb484b92..12b2f571 100644 --- a/internal/grpc/grpc_engine_test.go +++ b/internal/grpc/grpc_engine_test.go @@ -2,20 +2,130 @@ package grpc import ( "context" + "errors" + "path/filepath" "testing" + + "github.com/GrayCodeAI/flux/conversation" + "github.com/GrayCodeAI/flux/provider/core" + "github.com/GrayCodeAI/flux/storage" ) -// TestEngineChatServiceContract verifies the constructor surface and the noop -// fallback. A full engine-backed round-trip is covered by server_grpc_test.go -// (build tag "grpc") and requires a store-backed conversation.Engine. +// TestEngineChatServiceContract verifies the constructor surface and the +// placeholder paths. The engine-backed path is exercised below without the +// "grpc" build tag; server_grpc_test.go (tag "grpc") covers the wire framing +// with a stub ChatService. func TestEngineChatServiceContract(t *testing.T) { if NewChatService() == nil { t.Fatal("NewChatService returned nil") } + if _, err := NewChatService().Chat(context.Background(), &ChatRequest{Message: "hi"}); !errors.Is(err, ErrUnimplemented) { + t.Fatalf("expected ErrUnimplemented from the placeholder service, got %v", err) + } if svc := NewEngineChatService(nil); svc == nil { t.Fatal("NewEngineChatService returned nil") } - if _, err := NewEngineChatService(nil).Chat(context.Background(), &ChatRequest{Message: "hi"}); err != ErrUnimplemented { + if _, err := NewEngineChatService(nil).Chat(context.Background(), &ChatRequest{Message: "hi"}); !errors.Is(err, ErrUnimplemented) { t.Fatalf("expected ErrUnimplemented for nil engine, got %v", err) } } + +// scriptedProvider streams a fixed event sequence and records the request. +type scriptedProvider struct { + events []core.FluxStreamEvent + gotMsg chan []core.FluxMessage + gotOpt chan core.ChatOptions +} + +func (p *scriptedProvider) Name() string { return "scripted" } +func (p *scriptedProvider) Ping(_ context.Context) error { return nil } + +func (p *scriptedProvider) Chat(context.Context, []core.FluxMessage, core.ChatOptions) (*core.FluxResponse, error) { + return nil, errors.New("scriptedProvider: Chat is not used by conversation.Engine") +} + +func (p *scriptedProvider) StreamChat(_ context.Context, messages []core.FluxMessage, opts core.ChatOptions) (*core.StreamResult, error) { + p.gotMsg <- messages + p.gotOpt <- opts + ch := make(chan core.FluxStreamEvent, len(p.events)) + for _, evt := range p.events { + ch <- evt + } + close(ch) + return &core.StreamResult{Events: ch}, nil +} + +func newScriptedEngine(t *testing.T, events ...core.FluxStreamEvent) (*conversation.Engine, *scriptedProvider) { + t.Helper() + store, err := storage.Open(filepath.Join(t.TempDir(), "grpc.db")) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = store.Close() }) + prov := &scriptedProvider{ + events: events, + gotMsg: make(chan []core.FluxMessage, 1), + gotOpt: make(chan core.ChatOptions, 1), + } + return conversation.New(store, prov), prov +} + +func TestEngineChatServiceAggregatesEngineStream(t *testing.T) { + engine, prov := newScriptedEngine( + t, + core.FluxStreamEvent{Type: "content", Content: "hel"}, + core.FluxStreamEvent{Type: "content", Content: "lo"}, + core.FluxStreamEvent{Type: "done", StopReason: "end_turn", Usage: &core.FluxUsage{CompletionTokens: 2}}, + ) + + resp, err := NewEngineChatService(engine).Chat(context.Background(), &ChatRequest{ + Model: "test-model", + SystemPrompt: "be brief", + Message: "hi", + MaxTokens: 64, + }) + if err != nil { + t.Fatalf("Chat: %v", err) + } + if resp.Content != "hello" { + t.Fatalf("Content = %q, want the concatenated deltas %q", resp.Content, "hello") + } + if resp.NodeID == "" { + t.Fatal("NodeID is empty; want the saved assistant node ID") + } + if resp.FinishReason != "stop" { + t.Fatalf("FinishReason = %q, want stop", resp.FinishReason) + } + + msgs := <-prov.gotMsg + if len(msgs) != 1 || msgs[0].Role != "user" || msgs[0].Content != "hi" { + t.Fatalf("provider received messages %+v, want one user message %q", msgs, "hi") + } + opts := <-prov.gotOpt + if opts.Model != "test-model" || opts.System != "be brief" || opts.MaxTokens != 64 { + t.Fatalf("provider received options model=%q system=%q max=%d; want the ChatRequest fields", opts.Model, opts.System, opts.MaxTokens) + } +} + +func TestEngineChatServiceReturnsStreamError(t *testing.T) { + engine, _ := newScriptedEngine( + t, + core.FluxStreamEvent{Type: "content", Content: "partial"}, + core.FluxStreamEvent{Type: "error", Error: "upstream failed"}, + ) + + resp, err := NewEngineChatService(engine).Chat(context.Background(), &ChatRequest{Message: "hi"}) + if err == nil || err.Error() != "upstream failed" { + t.Fatalf("Chat error = %v, want the stream error %q", err, "upstream failed") + } + if resp != nil { + t.Fatalf("Chat response = %+v, want nil on stream error", resp) + } +} + +func TestEngineChatServiceRejectsNilRequest(t *testing.T) { + engine, _ := newScriptedEngine(t) + if _, err := NewEngineChatService(engine).Chat(context.Background(), nil); err == nil { + t.Fatal("Chat(nil) succeeded; want an error") + } +} diff --git a/runtime/runtime.go b/runtime/runtime.go index e180e475..7ae62f6d 100644 --- a/runtime/runtime.go +++ b/runtime/runtime.go @@ -1,25 +1,18 @@ -// Package runtime is the **recommended entry point** for host applications -// (e.g. rho). Start by calling runtime.Load to get a *Runtime, then -// rt.ChatProvider to obtain a core.Provider that you can hand to your -// agent loop. +// Package runtime resolves the active provider, model, deployment routing and +// credentials into a ready-to-use core.Provider: Load returns a *Runtime whose +// ChatProvider builds the transport for the current selection. // -// Note: the "stable" surface of flux is actually a set of cooperating -// subpackages, not just this one. The full list rho (and other host -// applications) actually import is: +// runtime is engine-internal. Hosts such as Rho must not import it; the host +// contract is limited to github.com/GrayCodeAI/flux/engine, llm, graph and +// tools (see docs/architecture/HOST-ENGINE-BOUNDARY.md), and engine calls this +// package on the host's behalf. The package stays importable for Flux's own +// packages and non-Rho integrations, but its exported names are not covered by +// engine.ContractVersion and may change in any release. // -// github.com/GrayCodeAI/flux/runtime (this package — bootstrap facade) -// github.com/GrayCodeAI/flux/provider (Provider interface, message/response types) -// github.com/GrayCodeAI/flux/catalog (model catalog: pricing, capabilities, registry) -// github.com/GrayCodeAI/flux/catalog/registry (ProviderSpec catalog: 16 registered providers) -// github.com/GrayCodeAI/flux/catalog/xiaomi (Xiaomi-specific catalog helpers) -// github.com/GrayCodeAI/flux/config (provider config + env var resolution) -// github.com/GrayCodeAI/flux/credentials (OS keyring + OIDC keyless CI auth) -// github.com/GrayCodeAI/flux/setup (CLI/setup wiring, RoutingPreviewJSON) -// github.com/GrayCodeAI/flux/storage (conversation DAG persistence) -// -// They are all considered part of the public API; changes to exported -// names are gated by semver. Anything under internal/ is implementation -// detail and may change without notice. +// Provider metadata (IDs, credential variables, protocols, regions) comes from +// the catalog registry in github.com/GrayCodeAI/flux/catalog/registry; do not +// restate its size here, it is checked against the docs by +// catalog/registry/docs_test.go. package runtime import ( diff --git a/scripts/test-config-flow.sh b/scripts/test-config-flow.sh index b0ccc00a..c455e635 100755 --- a/scripts/test-config-flow.sh +++ b/scripts/test-config-flow.sh @@ -1,7 +1,9 @@ #!/usr/bin/env bash # E2E test: /config flow — hub → credential → discover → picker → chat -# Run from flux root: bash scripts/test-config-flow.sh +# Run from anywhere: bash scripts/test-config-flow.sh +# Counts are derived from catalog/registry/providers.go, never hardcoded. set -euo pipefail +cd "$(dirname "$0")/.." PASS=0 FAIL=0 @@ -12,13 +14,13 @@ fail() { FAIL=$((FAIL+1)); echo " FAIL: $1"; } echo "=== Config Flow E2E Test ===" echo -# 1. Verify provider registry has all 11 providers +# 1. Verify the provider registry is populated echo "--- provider registry ---" -count=$(cd .. && grep -c "ProviderID:" flux/catalog/registry/providers.go 2>/dev/null || echo 0) -if [ "$count" -ge 11 ]; then +count=$(grep -c "ProviderID:" catalog/registry/providers.go || true) +if [ "${count:-0}" -gt 0 ]; then pass "registry has $count provider specs" else - fail "expected >= 11 providers, got $count" + fail "no ProviderID entries found in catalog/registry/providers.go" fi # 2. Verify all providers have deployment env fallbacks @@ -39,14 +41,13 @@ else fail "credential registry function not found" fi -# 4. Verify all providers have live fetchers +# 4. Verify every registry provider has a live fetcher echo "--- live fetchers ---" -cd "$(dirname "$0")/.." -fetchers=$(grep -c '".*":\s*Fetch' catalog/live/fetchers.go 2>/dev/null || echo 0) -if [ "$fetchers" -ge 11 ]; then - pass "all 11 providers have live fetchers" +fetchers=$(grep -cE '^[[:space:]]+"[a-z0-9_]+":[[:space:]]+Fetch' catalog/live/fetchers.go || true) +if [ "${fetchers:-0}" -eq "${count:-0}" ]; then + pass "all $count registry providers have live fetchers" else - fail "expected >= 11 fetchers, got $fetchers" + fail "registry has ${count:-0} providers but catalog/live/fetchers.go registers ${fetchers:-0} fetchers" fi # 5. Verify build + tests pass