adapters

package
v0.5.0 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Jul 25, 2026 License: Apache-2.0 Imports: 62 Imported by: 0

Documentation

Overview

Package adapters collects the convenience constructors for the concrete adapters shipped with agentflow: run-state/blob/memory stores, job queues, LLM providers, catalog manifests, knowledge and MCP tooling, built-in tool executors, tiered memory, and observability sinks/stores. It is a pure wiring layer: it never imports the agentflow root facade, so applications that only need these constructors can depend on it alone.

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

func LoadSkillManifest

func LoadSkillManifest(data []byte) (core.Skill, error)

LoadSkillManifest loads and validates a standalone skill catalog manifest document.

func LoadSkillManifestFile

func LoadSkillManifestFile(path string) (core.Skill, error)

LoadSkillManifestFile loads and validates a standalone skill catalog manifest.

func LoadToolManifest

func LoadToolManifest(data []byte) (core.Tool, error)

LoadToolManifest loads and validates a standalone tool catalog manifest document.

func LoadToolManifestFile

func LoadToolManifestFile(path string) (core.Tool, error)

LoadToolManifestFile loads and validates a standalone tool catalog manifest.

func NewAnthropicGateway

func NewAnthropicGateway(profiles []llm.Profile, client *http.Client) llm.Gateway

NewAnthropicGateway creates a gateway for Anthropic Messages APIs.

func NewBlobTierColdStore

func NewBlobTierColdStore(config BlobTierColdStoreConfig) (tier.Store, error)

NewBlobTierColdStore stores gzip JSON cold-tier records in a BlobAdmin backend.

func NewCognitiveTierMemory

func NewCognitiveTierMemory(manager tier.Manager, weights tier.RecallWeights) memory.CognitiveMemory

NewCognitiveTierMemory exposes a tier Manager through the CognitiveMemory port.

func NewCompositeTierStore

func NewCompositeTierStore(config CompositeTierStoreConfig) tier.Store

NewCompositeTierStore routes records across tier backends.

func NewEventFanoutSink

func NewEventFanoutSink(sinks ...core.EventSink) core.EventSink

func NewEventHub

func NewEventHub() *observability.EventHub

func NewEventStoreSink

func NewEventStoreSink(store observability.EventStore, publishers ...observability.EventPublisher) core.EventSink

func NewFileAuditSink

func NewFileAuditSink(path string) (audit.Sink, error)

func NewFileBlobStore

func NewFileBlobStore(dir string) (runstate.BlobStore, error)

NewFileBlobStore creates a file-backed blob store.

func NewFileKnowledgeLoader

func NewFileKnowledgeLoader(config FileKnowledgeLoaderConfig) (knowledge.Loader, error)

NewFileKnowledgeLoader creates a filesystem document loader for knowledge ingestion.

func NewFileMemoryRepository

func NewFileMemoryRepository(dir string) (memory.Repository, error)

NewFileMemoryRepository creates a JSON-file-backed memory repository.

func NewFileRunStateRepository

func NewFileRunStateRepository(dir string) (runstate.Repository, error)

NewFileRunStateRepository creates a JSON-file-backed run-state repository.

func NewFileTierColdStore

func NewFileTierColdStore(dir string) (tier.Store, error)

NewFileTierColdStore creates a gzip JSON cold-tier store on the local filesystem.

func NewFilesystemToolExecutor

func NewFilesystemToolExecutor(config FilesystemToolConfig) (core.ToolExecutor, error)

NewFilesystemToolExecutor creates a governed filesystem read tool executor.

func NewGitToolExecutor

func NewGitToolExecutor(config GitToolConfig) (core.ToolExecutor, error)

NewGitToolExecutor creates a read-only git tool executor.

func NewHTTPKnowledgeLoader

func NewHTTPKnowledgeLoader(config HTTPKnowledgeLoaderConfig) (knowledge.Loader, error)

NewHTTPKnowledgeLoader creates an HTTP document loader for knowledge ingestion.

func NewHTTPToolExecutor

func NewHTTPToolExecutor(config HTTPToolConfig) (core.ToolExecutor, error)

NewHTTPToolExecutor creates a governed HTTP client tool executor.

func NewInMemoryAuditSink

func NewInMemoryAuditSink(limit int) audit.Sink

func NewInMemoryBlobStore

func NewInMemoryBlobStore() runstate.BlobStore

NewInMemoryBlobStore creates the default in-memory blob store used by New.

func NewInMemoryCheckpointHistory

func NewInMemoryCheckpointHistory() runstate.CheckpointHistory

NewInMemoryCheckpointHistory creates an append-only in-memory checkpoint history store.

func NewInMemoryEventStore

func NewInMemoryEventStore() observability.EventStore

func NewInMemoryJobQueue

func NewInMemoryJobQueue() asyncpkg.Queue

func NewInMemoryRunStateRepository

func NewInMemoryRunStateRepository() runstate.Repository

NewInMemoryRunStateRepository creates the default in-memory run-state repository used by New.

func NewInMemoryTierHotStore

func NewInMemoryTierHotStore() tier.Store

NewInMemoryTierHotStore creates an in-process hot-tier store.

func NewKnowledgeIndexer

func NewKnowledgeIndexer(config KnowledgeIndexerConfig) (*knowledge.Indexer, error)

NewKnowledgeIndexer creates a document chunking, embedding, and vector upsert pipeline.

func NewLLMReranker

func NewLLMReranker(gateway llm.Gateway, profile string) knowledge.Reranker

NewLLMReranker creates an LLM reranker for retrieval tools.

func NewLLMRouter

func NewLLMRouter(routes map[string]llm.Gateway) llm.Gateway

NewLLMRouter routes profile names to provider-specific gateways.

func NewLLMTierSummarizer

func NewLLMTierSummarizer(gateway llm.Chatter, profile string) tier.ContentSummarizer

NewLLMTierSummarizer creates an LLM-backed cold-tier content summarizer.

func NewLocalGateway

func NewLocalGateway(profiles []llm.Profile, client *http.Client) llm.Gateway

NewLocalGateway creates a gateway for local OpenAI-compatible model servers.

func NewMCPHTTPClient

func NewMCPHTTPClient(endpoint string, client *http.Client) (mcp.Client, error)

NewMCPHTTPClient creates an MCP JSON-RPC client over HTTP.

func NewMCPHTTPClientWithOptions added in v0.4.4

func NewMCPHTTPClientWithOptions(endpoint string, client *http.Client, options mcp.ClientOptions) (mcp.Client, error)

NewMCPHTTPClientWithOptions creates an MCP HTTP client in the explicitly selected protocol mode.

func NewMCPToolExecutor

func NewMCPToolExecutor(client mcp.Client, tool string) (core.ToolExecutor, error)

NewMCPToolExecutor adapts one MCP server tool into an AgentFlow tool executor.

func NewNoopAuditSink

func NewNoopAuditSink() audit.Sink

func NewObservabilityEventSink

func NewObservabilityEventSink(recorder observability.Recorder, tracer observability.Tracer, next core.EventSink) core.EventSink

func NewOpenAICompatibleEmbedder

func NewOpenAICompatibleEmbedder(profiles []llm.Profile, client *http.Client) llm.Embedder

NewOpenAICompatibleEmbedder creates an embedder for OpenAI-compatible embedding APIs.

func NewOpenAICompatibleGateway

func NewOpenAICompatibleGateway(profiles []llm.Profile, client *http.Client) llm.Gateway

NewOpenAICompatibleGateway creates a gateway for OpenAI-compatible chat APIs.

func NewOpenTelemetryStdoutTracerProvider

func NewOpenTelemetryStdoutTracerProvider(ctx context.Context, config OpenTelemetryTracerProviderConfig) (*sdktrace.TracerProvider, error)

NewOpenTelemetryStdoutTracerProvider creates a TracerProvider that exports spans to stdout.

func NewOpenTelemetryTracer

func NewOpenTelemetryTracer(tracer oteltrace.Tracer) observability.Tracer

NewOpenTelemetryTracer wraps a host-configured OpenTelemetry tracer.

func NewPostgresCheckpointHistory

func NewPostgresCheckpointHistory(db *sql.DB, tableName ...string) (runstate.CheckpointHistory, error)

NewPostgresCheckpointHistory creates a PostgreSQL append-only checkpoint history store.

func NewPostgresJobQueue

func NewPostgresJobQueue(db *sql.DB, tableName ...string) (asyncpkg.Queue, error)

func NewPostgresOutboxEventSink added in v0.3.1

func NewPostgresOutboxEventSink(ctx context.Context, config PostgresOutboxEventSinkConfig, publishers ...observability.EventPublisher) (observability.EventStore, core.EventSink, error)

NewPostgresOutboxEventSink creates the PostgreSQL event store and an outbox-backed event sink for it. The sink appends to the store first; when the durable append fails, the event is parked in the run-state outbox (single INSERT, same database as the run snapshots) and later redelivered by the framework's outbox relay with its minted sequence, so a transient store outage no longer loses lifecycle events. Live publishers (e.g. an EventHub) are notified exactly like with NewEventStoreSink on the success path.

Wire the returned store into the framework with WithEventStore and start the relay with WithOutboxRelay, otherwise parked rows are never delivered:

store, sink, err := adapters.NewPostgresOutboxEventSink(ctx, cfg, eventHub)
fw, err := agentflow.New(scenario,
	agentflow.WithEventSink(adapters.NewEventFanoutSink(sink, ...)),
	agentflow.WithEventStore(store),
	agentflow.WithOutboxRelay(0),
	agentflow.WithRunStateRepository(pgRuns),
)

func NewPostgresRunStateRepository

func NewPostgresRunStateRepository(db *sql.DB, tableName ...string) (runstate.Repository, error)

NewPostgresRunStateRepository creates a PostgreSQL-compatible run-state repository using a caller-provided *sql.DB. Applications must import and register their preferred PostgreSQL database/sql driver.

func NewPostgresTierWarmStore

func NewPostgresTierWarmStore(config PostgresTierWarmStoreConfig) (tier.Store, error)

NewPostgresTierWarmStore creates a warm-tier store backed by Postgres JSONB rows.

func NewPostgresVectorStore

func NewPostgresVectorStore(config PostgresVectorStoreConfig) (knowledge.VectorStore, error)

NewPostgresVectorStore creates a pgvector-compatible knowledge vector store.

func NewRedisRunStateRepository

func NewRedisRunStateRepository(config RedisRunStateRepositoryConfig) (runstate.Repository, error)

NewRedisRunStateRepository creates a Redis-backed run-state repository with compare-and-swap version checks for distributed workers.

func NewRetrieverTool

func NewRetrieverTool(config RetrieverToolConfig) (core.ToolExecutor, error)

NewRetrieverTool creates a semantic retrieval tool backed by an embedder and vector store.

func NewS3BlobStore

func NewS3BlobStore(config S3BlobStoreConfig) (runstate.BlobStore, error)

NewS3BlobStore creates an S3-compatible blob store for large runtime and workflow outputs. It uses path-style object URLs, AWS Signature Version 4, and supports providers whose S3-compatible PUT/GET behavior has been tested.

func NewSQLToolExecutor

func NewSQLToolExecutor(config SQLToolConfig) (core.ToolExecutor, error)

NewSQLToolExecutor creates a governed read-only SQL query tool executor.

func NewScoreReranker

func NewScoreReranker() knowledge.Reranker

NewScoreReranker creates a lexical score reranker for retrieval tools.

func NewSlogAuditSink

func NewSlogAuditSink(logger *stdslog.Logger) audit.Sink

func NewSlogEventSink

func NewSlogEventSink(logger *stdslog.Logger) core.EventSink

func NewTicketToolExecutor

func NewTicketToolExecutor(config TicketToolConfig) (core.ToolExecutor, error)

NewTicketToolExecutor creates a ticket store backed tool executor.

func NewTierColdSummaryIndexer

func NewTierColdSummaryIndexer(config TierColdSummaryIndexerConfig) (tier.ColdSummaryIndexer, error)

NewTierColdSummaryIndexer indexes cold-tier summaries for semantic recall.

func NewVerboseSlogEventSink

func NewVerboseSlogEventSink(logger *stdslog.Logger) core.EventSink

NewVerboseSlogEventSink logs runtime events with redacted-safe payload details to stderr-friendly sinks.

func OpenTelemetryTracerFromProvider

func OpenTelemetryTracerFromProvider(provider *sdktrace.TracerProvider, instrumentationName string) observability.Tracer

OpenTelemetryTracerFromProvider returns a tracer backed by a TracerProvider.

func PrometheusMetricsHandler

func PrometheusMetricsHandler(recorder *PrometheusRecorder) http.Handler

PrometheusMetricsHandler returns an http.Handler that serves recorder metrics.

func ValidateSkillManifest

func ValidateSkillManifest(skill core.Skill) error

ValidateSkillManifest validates a skill manifest for catalog registration.

func ValidateToolManifest

func ValidateToolManifest(tool core.Tool) error

ValidateToolManifest validates a tool manifest for catalog registration.

Types

type BlobTierColdStoreConfig

type BlobTierColdStoreConfig struct {
	Blobs    runstate.BlobStore
	IndexDir string
}

BlobTierColdStoreConfig configures a blob-backed cold tier with a local index directory.

type CompositeTierStoreConfig

type CompositeTierStoreConfig struct {
	Hot  tier.Store
	Warm tier.Store
	Cold tier.Store
}

CompositeTierStoreConfig wires hot, warm, and cold tier backends.

type FileKnowledgeLoaderConfig

type FileKnowledgeLoaderConfig struct {
	Paths     []string
	Namespace string
	Metadata  map[string]string
	MaxBytes  int64
}

type FilesystemToolConfig

type FilesystemToolConfig struct {
	AllowedRoots []string
	MaxBytes     int64
}

type GitToolConfig

type GitToolConfig struct {
	AllowedRoots []string
}

type HTTPKnowledgeLoaderConfig

type HTTPKnowledgeLoaderConfig struct {
	URLs      []string
	Namespace string
	Metadata  map[string]string
	MaxBytes  int64
	Client    *http.Client
}

type HTTPToolConfig

type HTTPToolConfig struct {
	AllowedHosts     []string
	AllowedMethods   []string
	DefaultHeaders   map[string]string
	MaxResponseBytes int64
	Client           *http.Client
}

type KnowledgeIndexerConfig

type KnowledgeIndexerConfig struct {
	Embedder  llm.Embedder
	Store     knowledge.VectorStore
	Profile   string
	Namespace string
	BatchSize int
	Chunker   knowledge.Chunker
}

type LLMProviderRouter

type LLMProviderRouter interface {
	llm.Gateway
	llm.Embedder
}

func NewLLMProviderRouter

func NewLLMProviderRouter(routes map[string]llm.Gateway) LLMProviderRouter

NewLLMProviderRouter routes chat/tool/structured/streaming and embedding calls by profile name when the selected route supports the requested capability.

type MCPStdioClient

type MCPStdioClient interface {
	mcp.Client
	Close() error
}

func NewMCPStdioClient

func NewMCPStdioClient(ctx context.Context, config MCPStdioClientConfig) (MCPStdioClient, error)

NewMCPStdioClient creates an MCP JSON-RPC client over a child process stdio transport.

type MCPStdioClientConfig

type MCPStdioClientConfig struct {
	Command string
	Args    []string
	Env     []string
	Dir     string
	Options mcp.ClientOptions
}

type OpenAICompatibleProvider

type OpenAICompatibleProvider interface {
	llm.Gateway
	llm.Embedder
}

func NewOpenAICompatibleProvider

func NewOpenAICompatibleProvider(profiles []llm.Profile, client *http.Client) OpenAICompatibleProvider

NewOpenAICompatibleProvider creates a gateway/embedder for OpenAI-compatible APIs.

type OpenTelemetryTracer

type OpenTelemetryTracer = oteladapter.Tracer

OpenTelemetryTracer adapts go.opentelemetry.io/otel/trace.Tracer to observability.Tracer.

type OpenTelemetryTracerProviderConfig

type OpenTelemetryTracerProviderConfig = oteladapter.TracerProviderConfig

OpenTelemetryTracerProviderConfig configures a stdout-exporting TracerProvider for local development.

type PostgresEventStoreConfig

type PostgresEventStoreConfig struct {
	DB              *sql.DB
	TableName       string
	SkipSchemaSetup bool
}

type PostgresOutboxEventSinkConfig added in v0.3.1

type PostgresOutboxEventSinkConfig struct {
	DB *sql.DB
	// EventsTableName overrides the durable event table; see
	// PostgresEventStoreConfig.TableName.
	EventsTableName string
	// OutboxTableName overrides the fallback outbox table (default
	// agentflow_outbox, migration 0005). It must match the run-state
	// repository's outbox table when customized.
	OutboxTableName string
	SkipSchemaSetup bool
}

PostgresOutboxEventSinkConfig configures NewPostgresOutboxEventSink.

type PostgresTierWarmStoreConfig

type PostgresTierWarmStoreConfig struct {
	DB        *sql.DB
	TableName string
}

PostgresTierWarmStoreConfig configures a Postgres-backed warm tier store.

type PostgresVectorStoreConfig

type PostgresVectorStoreConfig struct {
	DB        *sql.DB
	TableName string
}

type PrometheusRecorder

type PrometheusRecorder = promrecorder.Recorder

PrometheusRecorder exposes in-process Prometheus text metrics for agentflow runtime signals.

func NewPrometheusRecorder

func NewPrometheusRecorder() *PrometheusRecorder

NewPrometheusRecorder creates a Prometheus-compatible observability recorder.

type RedisRunStateRepositoryConfig

type RedisRunStateRepositoryConfig struct {
	Addr         string
	Password     string
	DB           int
	KeyPrefix    string
	DialTimeout  time.Duration
	ReadTimeout  time.Duration
	WriteTimeout time.Duration
}

type RetrieverToolConfig

type RetrieverToolConfig struct {
	Embedder            llm.Embedder
	Store               knowledge.VectorStore
	Profile             string
	Namespace           string
	DefaultLimit        int
	SearchMode          knowledge.SearchMode
	CandidateMultiplier int
	Reranker            knowledge.Reranker
	VectorWeight        float64
	TextWeight          float64
}

type S3BlobStoreConfig

type S3BlobStoreConfig struct {
	Endpoint        string
	Bucket          string
	Region          string
	Prefix          string
	AccessKeyID     string
	SecretAccessKey string
	SessionToken    string
	HTTPClient      *http.Client
}

type SQLToolConfig

type SQLToolConfig struct {
	DB              *sql.DB
	AllowedQueries  map[string]string
	AllowAdHocQuery bool
	MaxRows         int
	Timeout         time.Duration
}

type Ticket

type Ticket = toolticket.Ticket

Ticket is a support ticket record manipulated by the ticket tool.

type TicketStore

type TicketStore = toolticket.Store

TicketStore persists ticket records for the ticket tool executor.

func NewMemoryTicketStore

func NewMemoryTicketStore(seed map[string]Ticket) TicketStore

NewMemoryTicketStore creates an in-memory ticket store for tests and demos.

type TicketToolConfig

type TicketToolConfig struct {
	Store toolticket.Store
}

type TierColdSummaryIndexerConfig

type TierColdSummaryIndexerConfig struct {
	Embedder   llm.Embedder
	Store      knowledge.VectorStore
	Profile    string
	MemoryName string
}

TierColdSummaryIndexerConfig configures vector indexing for cold-tier summaries.

type ToolResolver

type ToolResolver = core.ToolResolver

type ToolResolverFunc

type ToolResolverFunc = core.ToolResolverFunc

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL