mirror of
https://github.com/netbirdio/netbird.git
synced 2026-07-19 23:11:29 +02:00
* [agent-network] Shared proto, OpenAPI schema, and generated types * [agent-network] Management: store, manager, synthesizer, policy engine, provider catalog, HTTP/gRPC API Adds the account-scoped agent-network module: provider/policy/budget CRUD and store, the reverse-proxy service synthesizer, policy selection + limit enforcement, the provider catalog (incl. Vertex AI and AWS Bedrock entries), and the management HTTP + proxy gRPC surfaces. * [management] Fix agent-network proxy-peer fan-out on affected-peer recompute The affected-peers resolver loaded only persisted reverse-proxy services, but agent-network services are synthesized on demand and never persisted. As a result the embedded proxy peer was never folded into the affected set when a client's group changed, so the proxy received no network-map update for a newly authorised client and rejected its handshake until a full resync (restart). loadProxyServices now merges the synthesized agent-network services (injected via a registration hook to avoid an import cycle), so proxy peers learn newly authorised clients immediately. * [proxy] Reverse-proxy middleware framework, chain, and request plumbing The per-target middleware chain (slots, dispatcher, mutation gate, metadata merger), body capture, access-log terminal sink, and the proxy wiring that builds + runs chains for synthesized agent-network services. * [proxy] LLM parsers, pricing, and builtin middlewares (OpenAI, Anthropic, Vertex AI, AWS Bedrock) Request/response parsers and SSE/event-stream metering, the embedded pricing table, and the builtin middleware set: request parser, router, policy limit-check/record, cost meter, guardrail, identity inject, response parser. Includes the path-routed providers — Google Vertex AI (keyfile:: service-account OAuth minting) and AWS Bedrock (bearer auth, invoke/converse/streaming, optional /bedrock prefix) — plus the Models allowlist and unmeterable-publisher deny. * [proxy] IPv6 in-place apply and TCP accept-loop hardening on netstack listeners * [agent-network] End-to-end test suite, module docs, and deployment preset * [agent-network] Fix codespell typos and exclude false positives - labelgen word pool: vermillion -> vermilion, racoon -> raccoon. - codespell ignore list: add flate (Go compress/flate package), recordin (a test-local identifier), and unparseable (a valid alternative spelling used consistently across identifiers + a metadata-value constant). * [management] Set LastSeen on injected proxy peer in realstack test (MySQL strict-mode) The injected embedded proxy peer had a PeerStatus with a zero LastSeen, which serializes to '0000-00-00' and is rejected by MySQL in strict mode (SQLite tolerates it). Set LastSeen to a valid time so SaveAccount succeeds on both engines. * [agent-network] Remove e2e shell-script suite from this branch The end-to-end shell scripts under scripts/e2e/ are maintained in a separate testing suite and are not part of this change set. * [agent-network] Polish module docs: remove internal review scaffolding, fix links, verify diagrams Strip PR-review framing, commit references, absolute paths, and stale internal references from the agent-network module docs; fix broken relative links; verify all diagrams against the current architecture. Remove the internal AI-reviewer prompt file. * [management] Refine session expiration handling to support 3-state encoding for SSO deadlines * [agent-network] Relocate agentnetwork package to internals/modules Move management/server/agentnetwork (and its catalog/, labelgen/, types/ subpackages) to management/internals/modules/agentnetwork, alongside the reverse-proxy module, and rewrite all importers. Pure relocation: package names, the synthesizer + affectedpeers registration hook, and store access (shared store.Store) are unchanged, so no import cycle is introduced (affectedpeers still depends only on the agentnetwork/types leaf). * [agent-network] Co-locate HTTP handlers in the module (RegisterEndpoints) Move the agent-network HTTP handlers from server/http/handlers/agentnetwork into the module at internals/modules/agentnetwork/handlers (package handlers) and rename the entrypoint AddEndpoints -> RegisterEndpoints, matching the reverse-proxy module convention. Wiring in http/handler.go updated accordingly.
197 lines
6.2 KiB
Go
197 lines
6.2 KiB
Go
package llm
|
|
|
|
import (
|
|
"encoding/json"
|
|
"fmt"
|
|
"strings"
|
|
)
|
|
|
|
// AnthropicParser implements the Parser interface for the Anthropic Messages
|
|
// and Completions APIs. Detection is substring-based to tolerate upstream
|
|
// path rewrites.
|
|
type AnthropicParser struct{}
|
|
|
|
var anthropicPathHints = []string{
|
|
"/v1/messages",
|
|
"/v1/complete",
|
|
}
|
|
|
|
// Provider returns ProviderAnthropic.
|
|
func (AnthropicParser) Provider() Provider { return ProviderAnthropic }
|
|
|
|
// ProviderName returns the stable label used for metrics and metadata.
|
|
func (AnthropicParser) ProviderName() string { return "anthropic" }
|
|
|
|
// DetectFromURL reports whether the given request path looks like an
|
|
// Anthropic API endpoint. The match is case-insensitive and substring-based.
|
|
func (AnthropicParser) DetectFromURL(path string) bool {
|
|
lower := strings.ToLower(path)
|
|
for _, hint := range anthropicPathHints {
|
|
if strings.Contains(lower, hint) {
|
|
return true
|
|
}
|
|
}
|
|
return false
|
|
}
|
|
|
|
type anthropicRequest struct {
|
|
Model string `json:"model"`
|
|
Stream *bool `json:"stream"`
|
|
System json.RawMessage `json:"system"`
|
|
Messages []anthropicMessage `json:"messages"`
|
|
// Legacy /v1/complete endpoint.
|
|
Prompt string `json:"prompt"`
|
|
}
|
|
|
|
type anthropicMessage struct {
|
|
Role string `json:"role"`
|
|
Content json.RawMessage `json:"content"`
|
|
}
|
|
|
|
// ParseRequest extracts the model name and streaming flag from an Anthropic
|
|
// request body. Unknown or missing fields leave the corresponding struct
|
|
// members zero-valued.
|
|
func (AnthropicParser) ParseRequest(body []byte) (RequestFacts, error) {
|
|
var req anthropicRequest
|
|
if err := json.Unmarshal(body, &req); err != nil {
|
|
return RequestFacts{}, fmt.Errorf("decode anthropic request: %w: %v", ErrMalformedRequest, err)
|
|
}
|
|
return RequestFacts{
|
|
Model: req.Model,
|
|
Stream: ptrDeref(req.Stream),
|
|
}, nil
|
|
}
|
|
|
|
type anthropicResponse struct {
|
|
Usage struct {
|
|
InputTokens int64 `json:"input_tokens"`
|
|
OutputTokens int64 `json:"output_tokens"`
|
|
// CacheReadInputTokens and CacheCreationInputTokens are
|
|
// ADDITIVE to InputTokens (not subset), each billed at its
|
|
// own rate by the cost meter. cache_read is the cheaper
|
|
// read-from-cache rate, cache_creation is the more
|
|
// expensive write-to-cache rate.
|
|
CacheReadInputTokens int64 `json:"cache_read_input_tokens"`
|
|
CacheCreationInputTokens int64 `json:"cache_creation_input_tokens"`
|
|
} `json:"usage"`
|
|
}
|
|
|
|
// ParseResponse decodes the non-streaming Anthropic response envelope. Status
|
|
// codes other than 200 are treated as non-LLM responses so the caller can
|
|
// skip cost accounting without aborting the request.
|
|
func (AnthropicParser) ParseResponse(status int, contentType string, body []byte) (Usage, error) {
|
|
if status != 200 {
|
|
return Usage{}, fmt.Errorf("anthropic status %d: %w", status, ErrNotLLMResponse)
|
|
}
|
|
if isEventStream(contentType) {
|
|
return Usage{}, ErrStreamingUnsupported
|
|
}
|
|
if !isJSON(contentType) {
|
|
return Usage{}, fmt.Errorf("anthropic content-type %q: %w", contentType, ErrNotLLMResponse)
|
|
}
|
|
|
|
var resp anthropicResponse
|
|
if err := json.Unmarshal(body, &resp); err != nil {
|
|
return Usage{}, fmt.Errorf("decode anthropic response: %w: %v", ErrMalformedResponse, err)
|
|
}
|
|
return Usage{
|
|
InputTokens: resp.Usage.InputTokens,
|
|
OutputTokens: resp.Usage.OutputTokens,
|
|
TotalTokens: resp.Usage.InputTokens + resp.Usage.OutputTokens + resp.Usage.CacheReadInputTokens + resp.Usage.CacheCreationInputTokens,
|
|
CachedInputTokens: resp.Usage.CacheReadInputTokens,
|
|
CacheCreationTokens: resp.Usage.CacheCreationInputTokens,
|
|
}, nil
|
|
}
|
|
|
|
// ExtractPrompt returns the user-visible prompt text from an Anthropic
|
|
// request body. Handles the Messages API (system + messages[]) and the
|
|
// legacy /v1/complete prompt string. Returns "" on any decode failure.
|
|
func (AnthropicParser) ExtractPrompt(body []byte) string {
|
|
var req anthropicRequest
|
|
if err := json.Unmarshal(body, &req); err != nil {
|
|
return ""
|
|
}
|
|
var b strings.Builder
|
|
if len(req.System) > 0 {
|
|
if s := decodeStringOrJoin(req.System); s != "" {
|
|
b.WriteString("system: ")
|
|
b.WriteString(s)
|
|
}
|
|
}
|
|
for _, m := range req.Messages {
|
|
if b.Len() > 0 {
|
|
b.WriteByte('\n')
|
|
}
|
|
if m.Role != "" {
|
|
b.WriteString(m.Role)
|
|
b.WriteString(": ")
|
|
}
|
|
b.WriteString(decodeStringOrJoin(m.Content))
|
|
}
|
|
if b.Len() == 0 && req.Prompt != "" {
|
|
b.WriteString(req.Prompt)
|
|
}
|
|
return b.String()
|
|
}
|
|
|
|
// ExtractSessionID is the body-side fallback for Anthropic. Claude Code's
|
|
// authoritative session marker is the X-Claude-Code-Session-Id request
|
|
// header (handled by the request-parser middleware); this only mines the
|
|
// optional metadata.user_id for an embedded "...session_<uuid>" marker.
|
|
// metadata.user_id on its own is a USER identifier, not a session, so the
|
|
// whole value is deliberately NOT used — returning it would mislabel every
|
|
// request from a user as one session. Returns "" when no session marker is
|
|
// present.
|
|
func (AnthropicParser) ExtractSessionID(body []byte) string {
|
|
var req struct {
|
|
Metadata struct {
|
|
UserID string `json:"user_id"`
|
|
} `json:"metadata"`
|
|
}
|
|
if err := json.Unmarshal(body, &req); err != nil {
|
|
return ""
|
|
}
|
|
if idx := strings.LastIndex(req.Metadata.UserID, "session_"); idx >= 0 {
|
|
if session := req.Metadata.UserID[idx+len("session_"):]; session != "" {
|
|
return session
|
|
}
|
|
}
|
|
return ""
|
|
}
|
|
|
|
type anthropicMessageResponse struct {
|
|
Content []struct {
|
|
Type string `json:"type"`
|
|
Text string `json:"text"`
|
|
} `json:"content"`
|
|
// Legacy /v1/complete response.
|
|
Completion string `json:"completion"`
|
|
}
|
|
|
|
// ExtractCompletion returns the assistant text from a non-streaming Anthropic
|
|
// Messages or Completions response. Returns "" when status/content-type
|
|
// indicate the body is not parseable or no text part is present.
|
|
func (AnthropicParser) ExtractCompletion(status int, contentType string, body []byte) string {
|
|
if status != 200 || isEventStream(contentType) || !isJSON(contentType) {
|
|
return ""
|
|
}
|
|
var resp anthropicMessageResponse
|
|
if err := json.Unmarshal(body, &resp); err != nil {
|
|
return ""
|
|
}
|
|
var b strings.Builder
|
|
for _, part := range resp.Content {
|
|
if part.Text == "" {
|
|
continue
|
|
}
|
|
if b.Len() > 0 {
|
|
b.WriteByte('\n')
|
|
}
|
|
b.WriteString(part.Text)
|
|
}
|
|
if b.Len() == 0 {
|
|
return resp.Completion
|
|
}
|
|
return b.String()
|
|
}
|