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 ¶
- func LoadSkillManifest(data []byte) (core.Skill, error)
- func LoadSkillManifestFile(path string) (core.Skill, error)
- func LoadToolManifest(data []byte) (core.Tool, error)
- func LoadToolManifestFile(path string) (core.Tool, error)
- func NewAnthropicGateway(profiles []llm.Profile, client *http.Client) llm.Gateway
- func NewBlobTierColdStore(config BlobTierColdStoreConfig) (tier.Store, error)
- func NewCognitiveTierMemory(manager tier.Manager, weights tier.RecallWeights) memory.CognitiveMemory
- func NewCompositeTierStore(config CompositeTierStoreConfig) tier.Store
- func NewEventFanoutSink(sinks ...core.EventSink) core.EventSink
- func NewEventHub() *observability.EventHub
- func NewEventStoreSink(store observability.EventStore, publishers ...observability.EventPublisher) core.EventSink
- func NewFileAuditSink(path string) (audit.Sink, error)
- func NewFileBlobStore(dir string) (runstate.BlobStore, error)
- func NewFileKnowledgeLoader(config FileKnowledgeLoaderConfig) (knowledge.Loader, error)
- func NewFileMemoryRepository(dir string) (memory.Repository, error)
- func NewFileRunStateRepository(dir string) (runstate.Repository, error)
- func NewFileTierColdStore(dir string) (tier.Store, error)
- func NewFilesystemToolExecutor(config FilesystemToolConfig) (core.ToolExecutor, error)
- func NewGitToolExecutor(config GitToolConfig) (core.ToolExecutor, error)
- func NewHTTPKnowledgeLoader(config HTTPKnowledgeLoaderConfig) (knowledge.Loader, error)
- func NewHTTPToolExecutor(config HTTPToolConfig) (core.ToolExecutor, error)
- func NewInMemoryAuditSink(limit int) audit.Sink
- func NewInMemoryBlobStore() runstate.BlobStore
- func NewInMemoryCheckpointHistory() runstate.CheckpointHistory
- func NewInMemoryEventStore() observability.EventStore
- func NewInMemoryJobQueue() asyncpkg.Queue
- func NewInMemoryRunStateRepository() runstate.Repository
- func NewInMemoryTierHotStore() tier.Store
- func NewKnowledgeIndexer(config KnowledgeIndexerConfig) (*knowledge.Indexer, error)
- func NewLLMReranker(gateway llm.Gateway, profile string) knowledge.Reranker
- func NewLLMRouter(routes map[string]llm.Gateway) llm.Gateway
- func NewLLMTierSummarizer(gateway llm.Chatter, profile string) tier.ContentSummarizer
- func NewLocalGateway(profiles []llm.Profile, client *http.Client) llm.Gateway
- func NewMCPHTTPClient(endpoint string, client *http.Client) (mcp.Client, error)
- func NewMCPHTTPClientWithOptions(endpoint string, client *http.Client, options mcp.ClientOptions) (mcp.Client, error)
- func NewMCPToolExecutor(client mcp.Client, tool string) (core.ToolExecutor, error)
- func NewNoopAuditSink() audit.Sink
- func NewObservabilityEventSink(recorder observability.Recorder, tracer observability.Tracer, ...) core.EventSink
- func NewOpenAICompatibleEmbedder(profiles []llm.Profile, client *http.Client) llm.Embedder
- func NewOpenAICompatibleGateway(profiles []llm.Profile, client *http.Client) llm.Gateway
- func NewOpenTelemetryStdoutTracerProvider(ctx context.Context, config OpenTelemetryTracerProviderConfig) (*sdktrace.TracerProvider, error)
- func NewOpenTelemetryTracer(tracer oteltrace.Tracer) observability.Tracer
- func NewPostgresCheckpointHistory(db *sql.DB, tableName ...string) (runstate.CheckpointHistory, error)
- func NewPostgresEventStore(ctx context.Context, config PostgresEventStoreConfig) (observability.EventStore, error)
- func NewPostgresJobQueue(db *sql.DB, tableName ...string) (asyncpkg.Queue, error)
- func NewPostgresOutboxEventSink(ctx context.Context, config PostgresOutboxEventSinkConfig, ...) (observability.EventStore, core.EventSink, error)
- func NewPostgresRunStateRepository(db *sql.DB, tableName ...string) (runstate.Repository, error)
- func NewPostgresTierWarmStore(config PostgresTierWarmStoreConfig) (tier.Store, error)
- func NewPostgresVectorStore(config PostgresVectorStoreConfig) (knowledge.VectorStore, error)
- func NewRedisRunStateRepository(config RedisRunStateRepositoryConfig) (runstate.Repository, error)
- func NewRetrieverTool(config RetrieverToolConfig) (core.ToolExecutor, error)
- func NewS3BlobStore(config S3BlobStoreConfig) (runstate.BlobStore, error)
- func NewSQLToolExecutor(config SQLToolConfig) (core.ToolExecutor, error)
- func NewScoreReranker() knowledge.Reranker
- func NewSlogAuditSink(logger *stdslog.Logger) audit.Sink
- func NewSlogEventSink(logger *stdslog.Logger) core.EventSink
- func NewTicketToolExecutor(config TicketToolConfig) (core.ToolExecutor, error)
- func NewTierColdSummaryIndexer(config TierColdSummaryIndexerConfig) (tier.ColdSummaryIndexer, error)
- func NewVerboseSlogEventSink(logger *stdslog.Logger) core.EventSink
- func OpenTelemetryTracerFromProvider(provider *sdktrace.TracerProvider, instrumentationName string) observability.Tracer
- func PrometheusMetricsHandler(recorder *PrometheusRecorder) http.Handler
- func ValidateSkillManifest(skill core.Skill) error
- func ValidateToolManifest(tool core.Tool) error
- type BlobTierColdStoreConfig
- type CompositeTierStoreConfig
- type FileKnowledgeLoaderConfig
- type FilesystemToolConfig
- type GitToolConfig
- type HTTPKnowledgeLoaderConfig
- type HTTPToolConfig
- type KnowledgeIndexerConfig
- type LLMProviderRouter
- type MCPStdioClient
- type MCPStdioClientConfig
- type OpenAICompatibleProvider
- type OpenTelemetryTracer
- type OpenTelemetryTracerProviderConfig
- type PostgresEventStoreConfig
- type PostgresOutboxEventSinkConfig
- type PostgresTierWarmStoreConfig
- type PostgresVectorStoreConfig
- type PrometheusRecorder
- type RedisRunStateRepositoryConfig
- type RetrieverToolConfig
- type S3BlobStoreConfig
- type SQLToolConfig
- type Ticket
- type TicketStore
- type TicketToolConfig
- type TierColdSummaryIndexerConfig
- type ToolResolver
- type ToolResolverFunc
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
func LoadSkillManifest ¶
LoadSkillManifest loads and validates a standalone skill catalog manifest document.
func LoadSkillManifestFile ¶
LoadSkillManifestFile loads and validates a standalone skill catalog manifest.
func LoadToolManifest ¶
LoadToolManifest loads and validates a standalone tool catalog manifest document.
func LoadToolManifestFile ¶
LoadToolManifestFile loads and validates a standalone tool catalog manifest.
func NewAnthropicGateway ¶
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 NewEventHub ¶
func NewEventHub() *observability.EventHub
func NewEventStoreSink ¶
func NewEventStoreSink(store observability.EventStore, publishers ...observability.EventPublisher) core.EventSink
func NewFileBlobStore ¶
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 ¶
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 NewInMemoryBlobStore ¶
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 NewInMemoryRunStateRepository ¶
func NewInMemoryRunStateRepository() runstate.Repository
NewInMemoryRunStateRepository creates the default in-memory run-state repository used by New.
func NewInMemoryTierHotStore ¶
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 ¶
NewLLMReranker creates an LLM reranker for retrieval tools.
func NewLLMRouter ¶
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 ¶
NewLocalGateway creates a gateway for local OpenAI-compatible model servers.
func NewMCPHTTPClient ¶
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 ¶
NewMCPToolExecutor adapts one MCP server tool into an AgentFlow tool executor.
func NewNoopAuditSink ¶
func NewObservabilityEventSink ¶
func NewObservabilityEventSink(recorder observability.Recorder, tracer observability.Tracer, next core.EventSink) core.EventSink
func NewOpenAICompatibleEmbedder ¶
NewOpenAICompatibleEmbedder creates an embedder for OpenAI-compatible embedding APIs.
func NewOpenAICompatibleGateway ¶
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 NewPostgresEventStore ¶
func NewPostgresEventStore(ctx context.Context, config PostgresEventStoreConfig) (observability.EventStore, error)
func NewPostgresJobQueue ¶
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 ¶
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 ¶
NewScoreReranker creates a lexical score reranker for retrieval tools.
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 ¶
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 ¶
ValidateSkillManifest validates a skill manifest for catalog registration.
func ValidateToolManifest ¶
ValidateToolManifest validates a tool manifest for catalog registration.
Types ¶
type BlobTierColdStoreConfig ¶
BlobTierColdStoreConfig configures a blob-backed cold tier with a local index directory.
type CompositeTierStoreConfig ¶
CompositeTierStoreConfig wires hot, warm, and cold tier backends.
type FilesystemToolConfig ¶
type GitToolConfig ¶
type GitToolConfig struct {
AllowedRoots []string
}
type HTTPToolConfig ¶
type KnowledgeIndexerConfig ¶
type LLMProviderRouter ¶
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 ¶
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 OpenAICompatibleProvider ¶
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 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 ¶
PostgresTierWarmStoreConfig configures a Postgres-backed warm tier store.
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 RetrieverToolConfig ¶
type S3BlobStoreConfig ¶
type SQLToolConfig ¶
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