Files
synapbus/internal/mcp/server.go
T
Algis DumbrisandClaude Opus 4.6 cda3365863 feat(016): agent marketplace MVP — capability manifests, auction channels, reputation ledger
Implements US1, US2, and US3 of spec 016 by layering a marketplace service
on top of existing primitives rather than reinventing them:

- Capability manifests (US2) reuse the wiki subsystem. Each agent publishes
  a per-agent article at slug "agent-<name>" and gets versioning, revision
  history, and FTS search for free.
- Auction channels (US1) reuse the existing auction channel type, swarm
  service, and task/bid store. post_auction / bid / award wrap post_task /
  bid_task / accept_bid and attach marketplace metadata (max_budget_tokens,
  domains, estimated_tokens, confidence, approach) in the task.requirements
  and bid.capabilities JSON blobs. Award converts the auction into a claim
  by DM'ing the winner at priority 8 with task_id metadata, so the existing
  claim/process/done lifecycle takes over with zero new machinery.
- Reputation ledger (US3) adds migration 018_agent_marketplace.sql with a
  new agent_reputation table keyed by (agent_name, domain). mark_task_done
  completes the task via the swarm service and writes one ledger row per
  declared domain using the reported actual_tokens and success_score.
  query_reputation returns a rolled-up summary plus recent entries for a
  given (agent, domain) pair — reputation is always a vector, never a
  global score (FR-013).

Also:
- Adds the "awarded" reaction type (FR-008) alongside existing approve/
  reject/in_progress/done/published. Migration 018 widens the reactions
  CHECK constraint via a table rebuild.
- 6 new actions added to the action registry (post_auction, bid, award,
  mark_task_done, read_skill_card, query_reputation) so the search tool
  can discover them and the execute tool can dispatch them.
- New internal/marketplace package (store.go + service.go).
- New internal/mcp/marketplace.go bridge handlers.
- New internal/mcp/marketplace_test.go covers the full auction lifecycle,
  capability manifest publish/read/update, self-bid rejection, non-auction
  channel rejection, and reputation summary aggregation.

Out of scope for MVP (deferred per spec prompt): US4 reflection loop,
tombstoning FR-020a/b, multi-owner quorums, auto-escalation on zero bids,
bootstrap exploration credit, epsilon-greedy selection, and the hard-stop
budget enforcement daemon (only soft recording of estimated vs actual is
included).

All existing tests pass; new marketplace tests pass.

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
2026-04-11 15:17:30 +03:00

242 lines
7.0 KiB
Go

package mcp
import (
"context"
"database/sql"
"fmt"
"log/slog"
"net/http"
"time"
mcplib "github.com/mark3labs/mcp-go/mcp"
"github.com/mark3labs/mcp-go/server"
"github.com/synapbus/synapbus/internal/actions"
"github.com/synapbus/synapbus/internal/agentquery"
"github.com/synapbus/synapbus/internal/agents"
"github.com/synapbus/synapbus/internal/attachments"
"github.com/synapbus/synapbus/internal/channels"
"github.com/synapbus/synapbus/internal/console"
"github.com/synapbus/synapbus/internal/jsruntime"
"github.com/synapbus/synapbus/internal/marketplace"
"github.com/synapbus/synapbus/internal/messaging"
"github.com/synapbus/synapbus/internal/reactions"
"github.com/synapbus/synapbus/internal/search"
"github.com/synapbus/synapbus/internal/trace"
"github.com/synapbus/synapbus/internal/trust"
"github.com/synapbus/synapbus/internal/wiki"
)
// MCPServer wraps the mcp-go server with SynapBus services.
type MCPServer struct {
mcpServer *server.MCPServer
httpServer *server.StreamableHTTPServer
connMgr *ConnectionManager
agentService *agents.AgentService
hybridRegistrar *HybridToolRegistrar
logger *slog.Logger
console *console.Printer
}
// NewMCPServer creates and configures a new MCP server with 5 hybrid tools registered.
func NewMCPServer(
msgService *messaging.MessagingService,
agentService *agents.AgentService,
channelService *channels.Service,
swarmService *channels.SwarmService,
attachmentService *attachments.Service,
searchService *search.Service,
reactionService *reactions.Service,
trustService *trust.Service,
wikiService *wiki.Service,
consolePrinter *console.Printer,
jsPool *jsruntime.Pool,
actionRegistry *actions.Registry,
actionIndex *actions.Index,
db *sql.DB,
) *MCPServer {
logger := slog.Default().With("component", "mcp-server")
connMgr := NewConnectionManager()
// Set up hooks for client info capture and connection tracking
hooks := &server.Hooks{}
hooks.AddAfterInitialize(func(ctx context.Context, id any, msg *mcplib.InitializeRequest, result *mcplib.InitializeResult) {
clientName := msg.Params.ClientInfo.Name
clientVersion := msg.Params.ClientInfo.Version
protocolVersion := msg.Params.ProtocolVersion
// Build capabilities list
var caps []string
if msg.Params.Capabilities.Roots != nil {
caps = append(caps, "roots")
}
if msg.Params.Capabilities.Sampling != nil {
caps = append(caps, "sampling")
}
if msg.Params.Capabilities.Elicitation != nil {
caps = append(caps, "elicitation")
}
if len(msg.Params.Capabilities.Experimental) > 0 {
for k := range msg.Params.Capabilities.Experimental {
caps = append(caps, "experimental/"+k)
}
}
// Extract agent name from context (set by HTTP auth middleware)
agentName, _ := extractAgentName(ctx)
// Get session ID for connection tracking
session := server.ClientSessionFromContext(ctx)
sessionID := ""
if session != nil {
sessionID = session.SessionID()
}
// Register connection
if sessionID != "" {
conn := &Connection{
ID: sessionID,
AgentName: agentName,
Transport: "streamable-http",
ConnectedAt: time.Now(),
LastActivity: time.Now(),
ClientName: clientName,
ClientVersion: clientVersion,
ProtocolVersion: protocolVersion,
ClientCapabilities: caps,
}
connMgr.Add(conn)
}
// Structured log (always)
logger.Info("client initialized",
"agent", agentName,
"client_name", clientName,
"client_version", clientVersion,
"protocol_version", protocolVersion,
"capabilities", caps,
"session_id", sessionID,
)
// Pretty console output
if consolePrinter != nil {
if agentName != "" {
consolePrinter.AgentConnected(agentName, clientName, clientVersion)
} else {
consolePrinter.ClientConnected(clientName, clientVersion)
}
}
})
hooks.AddOnUnregisterSession(func(ctx context.Context, session server.ClientSession) {
sessionID := session.SessionID()
conn, ok := connMgr.Get(sessionID)
if ok {
if consolePrinter != nil && conn.AgentName != "" {
consolePrinter.AgentDisconnected(conn.AgentName)
}
logger.Info("client disconnected",
"agent", conn.AgentName,
"client_name", conn.ClientName,
"session_id", sessionID,
"duration", fmt.Sprintf("%s", time.Since(conn.ConnectedAt).Truncate(time.Second)),
)
connMgr.Remove(sessionID)
}
})
// Create the mcp-go server
mcpSrv := server.NewMCPServer(
"SynapBus",
"0.1.0",
server.WithToolCapabilities(true),
server.WithPromptCapabilities(true),
server.WithHooks(hooks),
)
// Register the 4 hybrid tools
hybridRegistrar := NewHybridToolRegistrar(
msgService,
agentService,
channelService,
swarmService,
attachmentService,
searchService,
reactionService,
trustService,
wikiService,
jsPool,
actionRegistry,
actionIndex,
db,
)
hybridRegistrar.RegisterAllOnServer(mcpSrv)
// Register the 4 MCP prompts
traceStore := trace.NewSQLiteTraceStore(db)
promptRegistrar := NewPromptRegistrar(db, agentService, channelService, traceStore)
promptRegistrar.RegisterAllOnServer(mcpSrv)
// Create Streamable HTTP transport with context func for auth propagation
httpServer := server.NewStreamableHTTPServer(mcpSrv,
server.WithHTTPContextFunc(func(ctx context.Context, r *http.Request) context.Context {
// Propagate agent identity from HTTP auth to MCP context
if agent, ok := agents.AgentFromContext(r.Context()); ok {
ctx = ContextWithAgentName(ctx, agent.Name)
// Propagate owner ID for trace recording
if ownerID, ok := trace.OwnerIDFromContext(r.Context()); ok {
ctx = trace.ContextWithOwnerID(ctx, ownerID)
}
}
return ctx
}),
)
s := &MCPServer{
mcpServer: mcpSrv,
httpServer: httpServer,
connMgr: connMgr,
agentService: agentService,
hybridRegistrar: hybridRegistrar,
logger: logger,
console: consolePrinter,
}
logger.Info("MCP server initialized (4 hybrid tools, 4 prompts, streamable HTTP transport)")
return s
}
// SetQueryExecutor sets the SQL query executor for agent queries via the execute tool.
func (s *MCPServer) SetQueryExecutor(exec *agentquery.Executor) {
if s.hybridRegistrar != nil {
s.hybridRegistrar.SetQueryExecutor(exec)
}
}
// SetMarketplaceService wires the marketplace service (spec 016) into the
// hybrid tool registrar so execute-tool calls can reach marketplace actions.
func (s *MCPServer) SetMarketplaceService(m *marketplace.Service) {
if s.hybridRegistrar != nil {
s.hybridRegistrar.SetMarketplaceService(m)
}
}
// Handler returns the HTTP handler for mounting on a router.
func (s *MCPServer) Handler() http.Handler {
return s.httpServer
}
// ConnectionManager returns the connection manager.
func (s *MCPServer) ConnectionManager() *ConnectionManager {
return s.connMgr
}
// Shutdown gracefully shuts down the MCP server.
func (s *MCPServer) Shutdown(ctx context.Context) error {
s.logger.Info("MCP server shutting down")
if err := s.httpServer.Shutdown(ctx); err != nil {
return err
}
return nil
}