From cda336586334d42cead71690f0db76911b7be8d4 Mon Sep 17 00:00:00 2001 From: Algis Dumbris Date: Sat, 11 Apr 2026 15:17:30 +0300 Subject: [PATCH] =?UTF-8?q?feat(016):=20agent=20marketplace=20MVP=20?= =?UTF-8?q?=E2=80=94=20capability=20manifests,=20auction=20channels,=20rep?= =?UTF-8?q?utation=20ledger?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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-" 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) --- cmd/synapbus/main.go | 7 + internal/actions/registry.go | 103 ++++- internal/actions/registry_test.go | 10 +- internal/marketplace/service.go | 356 ++++++++++++++++ internal/marketplace/store.go | 196 +++++++++ internal/mcp/bridge.go | 16 + internal/mcp/marketplace.go | 296 ++++++++++++++ internal/mcp/marketplace_test.go | 379 ++++++++++++++++++ internal/mcp/server.go | 9 + internal/mcp/tools_hybrid.go | 12 + internal/reactions/model.go | 13 +- .../storage/schema/018_agent_marketplace.sql | 51 +++ 12 files changed, 1440 insertions(+), 8 deletions(-) create mode 100644 internal/marketplace/service.go create mode 100644 internal/marketplace/store.go create mode 100644 internal/mcp/marketplace.go create mode 100644 internal/mcp/marketplace_test.go create mode 100644 internal/storage/schema/018_agent_marketplace.sql diff --git a/cmd/synapbus/main.go b/cmd/synapbus/main.go index af36693..8ed9718 100644 --- a/cmd/synapbus/main.go +++ b/cmd/synapbus/main.go @@ -38,6 +38,7 @@ import ( "github.com/synapbus/synapbus/internal/health" "github.com/synapbus/synapbus/internal/jsruntime" k8spkg "github.com/synapbus/synapbus/internal/k8s" + "github.com/synapbus/synapbus/internal/marketplace" mcpserver "github.com/synapbus/synapbus/internal/mcp" "github.com/synapbus/synapbus/internal/agentquery" reactorpkg "github.com/synapbus/synapbus/internal/reactor" @@ -500,6 +501,12 @@ func runServe(cmd *cobra.Command, args []string) error { mcpSrv := mcpserver.NewMCPServer(msgService, agentService, channelService, swarmService, attachmentService, searchService, reactionService, trustService, wikiService, con, jsPool, actionRegistry, actionIndex, db.DB) + // Wire the agent marketplace (spec 016 MVP). + marketplaceStore := marketplace.NewStore(db.DB) + marketplaceSvc := marketplace.NewService(marketplaceStore, wikiService, swarmService, channelService, msgService, tracer) + mcpSrv.SetMarketplaceService(marketplaceSvc) + slog.Info("agent marketplace service initialized (spec 016)") + // Set up SQL query executor for agents (uses read pool if available) queryDB := db.QueryDB() queryExec := agentquery.New(queryDB, slog.Default()) diff --git a/internal/actions/registry.go b/internal/actions/registry.go index c4e1c89..f24d3a3 100644 --- a/internal/actions/registry.go +++ b/internal/actions/registry.go @@ -6,10 +6,12 @@ type Registry struct { ordered []Action // maintains insertion order } -// NewRegistry creates a registry pre-populated with all 33 agent-callable actions. +// NewRegistry creates a registry pre-populated with the full set of +// agent-callable actions (marketplace additions in spec 016 bring the +// total to ~39). func NewRegistry() *Registry { r := &Registry{ - actions: make(map[string]Action, 33), + actions: make(map[string]Action, 40), } for _, a := range allActions() { r.actions[a.Name] = a @@ -672,5 +674,102 @@ func allActions() []Action { Code: `call("get_backlinks", {"slug": "mcp-security-landscape"})`, }}, }, + + // ── Marketplace (6 actions) — spec 016 ────────────────────── + { + Name: "post_auction", + Category: "marketplace", + Description: "Post an auction task to an auction-type channel (spec 016 / US1). Agents bid on the task and the poster awards one bid. Include max_budget_tokens and at least one domain tag so reputation can be scoped when the task completes.", + Params: []Param{ + {Name: "channel_name", Type: "string", Required: true, Description: "Name of the auction channel"}, + {Name: "title", Type: "string", Required: true, Description: "Short task title"}, + {Name: "description", Type: "string", Description: "Task description"}, + {Name: "acceptance_criteria", Type: "string", Description: "What success looks like"}, + {Name: "max_budget_tokens", Type: "number", Description: "Maximum token budget for the winning agent"}, + {Name: "domains", Type: "string", Description: "Comma-separated domain tags (e.g. 'data-analysis,python') or JSON array"}, + {Name: "difficulty_weight", Type: "number", Description: "Difficulty multiplier for reputation scoring (default 1.0)"}, + {Name: "deadline", Type: "string", Description: "ISO 8601 deadline"}, + }, + Returns: "JSON with task_id, channel_id, title, status, domains, max_budget_tokens, deadline", + Examples: []Example{{ + Description: "Post an auction for a data analysis task", + Code: `call("post_auction", {"channel_name": "task-marketplace", "title": "Q4 revenue analysis", "description": "Trend analysis with charts", "domains": "data-analysis,python", "max_budget_tokens": 8000, "deadline": "2026-05-01T17:00:00Z"})`, + }}, + }, + { + Name: "bid", + Category: "marketplace", + Description: "Submit a bid on an auction task (spec 016 / US1). Include estimated_tokens, a confidence score in [0,1], a brief approach summary, and the revision of your capability manifest at time of bid.", + Params: []Param{ + {Name: "task_id", Type: "number", Required: true, Description: "ID of the auction task"}, + {Name: "estimated_tokens", Type: "number", Description: "Your estimated token cost to complete the task"}, + {Name: "confidence", Type: "number", Description: "Self-reported confidence in 0.0..1.0"}, + {Name: "approach", Type: "string", Description: "Brief approach summary"}, + {Name: "manifest_revision", Type: "number", Description: "Revision of your capability manifest at time of bid"}, + {Name: "time_estimate", Type: "string", Description: "Optional human-readable time estimate"}, + }, + Returns: "JSON with bid_id, task_id, agent_name, status, estimated_tokens, confidence", + Examples: []Example{{ + Description: "Bid on task 42 with an 8k token estimate", + Code: `call("bid", {"task_id": 42, "estimated_tokens": 4200, "confidence": 0.9, "approach": "Pandas + matplotlib"})`, + }}, + }, + { + Name: "award", + Category: "marketplace", + Description: "Award an auction task to a specific bid (spec 016 / US1). Only the task poster can award. On award, the winning agent receives a high-priority DM they can process via the normal claim/process/done lifecycle.", + Params: []Param{ + {Name: "task_id", Type: "number", Required: true, Description: "ID of the task"}, + {Name: "bid_id", Type: "number", Required: true, Description: "ID of the winning bid"}, + }, + Returns: "JSON with task_id, bid_id, winner, claim_message_id, status", + Examples: []Example{{ + Description: "Award task 42 to bid 7", + Code: `call("award", {"task_id": 42, "bid_id": 7})`, + }}, + }, + { + Name: "mark_task_done", + Category: "marketplace", + Description: "Mark an auction task done (spec 016 / US3). Only the assigned agent can call this. Records reputation ledger entries for each declared domain using the reported actual_tokens and success_score.", + Params: []Param{ + {Name: "task_id", Type: "number", Required: true, Description: "ID of the assigned task"}, + {Name: "actual_tokens", Type: "number", Description: "Actual tokens spent"}, + {Name: "success_score", Type: "number", Description: "Self-reported success score in 0.0..1.0 (default 1.0)"}, + }, + Returns: "JSON with task_id, status, actual_tokens, success_score, reputation_entries", + Examples: []Example{{ + Description: "Mark task 42 completed with 3800 tokens spent", + Code: `call("mark_task_done", {"task_id": 42, "actual_tokens": 3800, "success_score": 1.0})`, + }}, + }, + { + Name: "read_skill_card", + Category: "marketplace", + Description: "Read an agent's capability manifest (spec 016 / US2). Capability manifests are versioned wiki articles at slug 'agent-'. Returns exists=false if the agent has not published a manifest yet. Use create_article / update_article on slug 'agent-' to publish or update your own.", + Params: []Param{ + {Name: "agent_name", Type: "string", Description: "Name of the agent whose manifest you want to read (defaults to caller)"}, + }, + Returns: "JSON with agent_name, exists, slug, title, body, revision, updated_at", + Examples: []Example{{ + Description: "Read another agent's skill card", + Code: `call("read_skill_card", {"agent_name": "data-processor"})`, + }}, + }, + { + Name: "query_reputation", + Category: "marketplace", + Description: "Query the per-(agent, domain) reputation ledger (spec 016 / US3). Reputation is a vector — you must supply both agent and domain. Returns an aggregated summary and the most recent raw ledger entries.", + Params: []Param{ + {Name: "agent_name", Type: "string", Description: "Agent to query (defaults to caller)"}, + {Name: "domain", Type: "string", Required: true, Description: "Domain tag to scope the query"}, + {Name: "limit", Type: "number", Description: "Max raw entries to return (default 20)"}, + }, + Returns: "JSON with agent_name, domain, summary (tasks_completed, avg_success_score, weighted_success_score, avg_estimated_tokens, avg_actual_tokens), recent_entries", + Examples: []Example{{ + Description: "Check an agent's reputation in data-analysis", + Code: `call("query_reputation", {"agent_name": "data-processor", "domain": "data-analysis"})`, + }}, + }, } } diff --git a/internal/actions/registry_test.go b/internal/actions/registry_test.go index 68a5b62..b56651c 100644 --- a/internal/actions/registry_test.go +++ b/internal/actions/registry_test.go @@ -4,11 +4,12 @@ import ( "testing" ) -func TestRegistryHas35Actions(t *testing.T) { +func TestRegistryHasAllActions(t *testing.T) { r := NewRegistry() got := len(r.List()) - if got != 35 { - t.Errorf("expected 35 actions, got %d", got) + const want = 41 // 35 original + 6 marketplace (spec 016) + if got != want { + t.Errorf("expected %d actions, got %d", want, got) } } @@ -28,6 +29,7 @@ func TestRegistryCategories(t *testing.T) { {"trust", 1}, {"data", 1}, {"wiki", 5}, + {"marketplace", 6}, } for _, tt := range tests { @@ -64,6 +66,8 @@ func TestRegistryGetByName(t *testing.T) { "query", // wiki "create_article", "get_article", "update_article", "list_articles", "get_backlinks", + // marketplace (spec 016) + "post_auction", "bid", "award", "mark_task_done", "read_skill_card", "query_reputation", } for _, name := range allNames { diff --git a/internal/marketplace/service.go b/internal/marketplace/service.go new file mode 100644 index 0000000..f0b9c41 --- /dev/null +++ b/internal/marketplace/service.go @@ -0,0 +1,356 @@ +package marketplace + +import ( + "context" + "encoding/json" + "fmt" + "log/slog" + "strings" + "time" + + "github.com/synapbus/synapbus/internal/channels" + "github.com/synapbus/synapbus/internal/messaging" + "github.com/synapbus/synapbus/internal/trace" + "github.com/synapbus/synapbus/internal/wiki" +) + +// SkillCardSlug returns the wiki slug for an agent's capability manifest. +// Wiki slugs allow only [a-z0-9-], so we use "agent-". Any underscores +// in the agent name are replaced with hyphens. +func SkillCardSlug(agentName string) string { + a := strings.ToLower(strings.TrimSpace(agentName)) + a = strings.ReplaceAll(a, "_", "-") + return "agent-" + a +} + +// Service implements the marketplace MVP (spec 016). +type Service struct { + store *Store + wikiService *wiki.Service + swarm *channels.SwarmService + channels *channels.Service + messaging *messaging.MessagingService + tracer *trace.Tracer + logger *slog.Logger +} + +// NewService wires a marketplace service. All collaborators are required +// except tracer (nil-safe). +func NewService( + store *Store, + wikiService *wiki.Service, + swarm *channels.SwarmService, + channelService *channels.Service, + msgService *messaging.MessagingService, + tracer *trace.Tracer, +) *Service { + return &Service{ + store: store, + wikiService: wikiService, + swarm: swarm, + channels: channelService, + messaging: msgService, + tracer: tracer, + logger: slog.Default().With("component", "marketplace"), + } +} + +// ---------- Types carried in task.Requirements / bid.Capabilities JSON. ---------- + +// AuctionTaskMeta is stored inside tasks.requirements as JSON. All fields +// are optional on the wire; missing values fall back to sane defaults. +type AuctionTaskMeta struct { + AcceptanceCriteria string `json:"acceptance_criteria,omitempty"` + MaxBudgetTokens int64 `json:"max_budget_tokens,omitempty"` + Domains []string `json:"domains,omitempty"` + DifficultyWeight float64 `json:"difficulty_weight,omitempty"` + CurrentSpendTokens int64 `json:"current_spend_tokens,omitempty"` +} + +// BidMeta is stored inside task_bids.capabilities as JSON. +type BidMeta struct { + EstimatedTokens int64 `json:"estimated_tokens,omitempty"` + Confidence float64 `json:"confidence,omitempty"` + Approach string `json:"approach,omitempty"` + ManifestRevision int `json:"manifest_revision,omitempty"` +} + +// ---------- Capability manifest (US2 via wiki). ---------- + +// ReadSkillCard returns the capability manifest article for the given agent. +// Returns a nil article + nil error if the manifest does not yet exist. +func (s *Service) ReadSkillCard(ctx context.Context, agentName string) (*wiki.Article, error) { + if s.wikiService == nil { + return nil, fmt.Errorf("wiki service not available") + } + slug := SkillCardSlug(agentName) + art, err := s.wikiService.GetArticle(ctx, slug) + if err != nil { + // Treat "not found" as a nil result — the caller decides whether this + // is an error. The wiki store returns an error string containing + // "not found" when the slug is missing. + if strings.Contains(err.Error(), "not found") { + return nil, nil + } + return nil, err + } + return art, nil +} + +// ---------- Auction posting / bidding / awarding (US1 on top of swarm). ---------- + +// PostAuction creates a new auction task in an auction-type channel with the +// supplied marketplace metadata serialised into task.requirements. +func (s *Service) PostAuction( + ctx context.Context, + channelID int64, + agentName, title, description string, + meta AuctionTaskMeta, + deadline *time.Time, +) (*channels.Task, error) { + if s.swarm == nil { + return nil, fmt.Errorf("swarm service not available") + } + + // Normalise meta — ensure domains slice is non-nil and difficulty defaults to 1. + if meta.Domains == nil { + meta.Domains = []string{} + } + if meta.DifficultyWeight == 0 { + meta.DifficultyWeight = 1.0 + } + + reqs, err := json.Marshal(meta) + if err != nil { + return nil, fmt.Errorf("marshal auction meta: %w", err) + } + + task, err := s.swarm.PostTask(ctx, channelID, agentName, title, description, reqs, deadline) + if err != nil { + return nil, err + } + + if s.tracer != nil { + s.tracer.Record(ctx, agentName, "marketplace.auction_posted", map[string]any{ + "task_id": task.ID, + "channel_id": channelID, + "domains": meta.Domains, + "max_budget_tokens": meta.MaxBudgetTokens, + }) + } + return task, nil +} + +// Bid submits a bid on an open auction task. Delegates to the swarm service +// for the core bid lifecycle and stores marketplace bid metadata in the +// bid.capabilities JSON blob. +func (s *Service) Bid( + ctx context.Context, + taskID int64, + agentName string, + meta BidMeta, + timeEstimate string, +) (*channels.Bid, error) { + if s.swarm == nil { + return nil, fmt.Errorf("swarm service not available") + } + + caps, err := json.Marshal(meta) + if err != nil { + return nil, fmt.Errorf("marshal bid meta: %w", err) + } + + // Use meta.Approach as the bid message for human readability. + bid, err := s.swarm.BidOnTask(ctx, taskID, agentName, caps, timeEstimate, meta.Approach) + if err != nil { + return nil, err + } + + if s.tracer != nil { + s.tracer.Record(ctx, agentName, "marketplace.bid_submitted", map[string]any{ + "task_id": taskID, + "bid_id": bid.ID, + "estimated_tokens": meta.EstimatedTokens, + "confidence": meta.Confidence, + }) + } + return bid, nil +} + +// Award accepts a bid and converts the auction into a claim on the winning +// agent by DM'ing them. The task body contains the task_id in metadata so the +// winning agent can use the existing claim/process/done lifecycle. +// +// Per spec 016 FR-009: "On award, the system MUST convert the auction into a +// claim on the winning agent using the existing claim/process/done lifecycle." +func (s *Service) Award( + ctx context.Context, + taskID, bidID int64, + awarderAgent string, +) (winningAgent string, claimMessageID int64, err error) { + if s.swarm == nil { + return "", 0, fmt.Errorf("swarm service not available") + } + + if err := s.swarm.AcceptBid(ctx, taskID, bidID, awarderAgent); err != nil { + return "", 0, err + } + + // Refresh task to find the assigned agent. + task, bids, err := s.swarm.GetTaskWithBids(ctx, taskID) + if err != nil { + return "", 0, fmt.Errorf("load awarded task: %w", err) + } + winningAgent = task.AssignedTo + if winningAgent == "" { + // Fallback: find the accepted bid. + for _, b := range bids { + if b.ID == bidID { + winningAgent = b.AgentName + break + } + } + } + if winningAgent == "" { + return "", 0, fmt.Errorf("could not determine winning agent for task %d", taskID) + } + + // Send a DM to the winner that acts as the claimable work item. + if s.messaging != nil { + metaJSON, _ := json.Marshal(map[string]any{ + "marketplace": "awarded", + "task_id": taskID, + "bid_id": bidID, + "channel_id": task.ChannelID, + "awarded_by": awarderAgent, + }) + body := fmt.Sprintf("Auction awarded: task %d — %q. Use mark_task_done to complete.", task.ID, task.Title) + msg, sendErr := s.messaging.SendMessage(ctx, awarderAgent, winningAgent, body, messaging.SendOptions{ + Subject: "Awarded: " + task.Title, + Priority: 8, + Metadata: string(metaJSON), + }) + if sendErr != nil { + s.logger.Warn("failed to send award DM", "err", sendErr, "task_id", taskID, "winner", winningAgent) + } else if msg != nil { + claimMessageID = msg.ID + } + } + + if s.tracer != nil { + s.tracer.Record(ctx, awarderAgent, "marketplace.auction_awarded", map[string]any{ + "task_id": taskID, + "bid_id": bidID, + "winner": winningAgent, + "claim_message_id": claimMessageID, + }) + } + return winningAgent, claimMessageID, nil +} + +// MarkTaskDone marks a task completed and writes reputation ledger entries +// for every domain declared on the task. One entry per domain (FR-012). +func (s *Service) MarkTaskDone( + ctx context.Context, + taskID int64, + agentName string, + actualTokens int64, + successScore float64, +) (*channels.Task, []*ReputationEntry, error) { + if s.swarm == nil { + return nil, nil, fmt.Errorf("swarm service not available") + } + + task, bids, err := s.swarm.GetTaskWithBids(ctx, taskID) + if err != nil { + return nil, nil, err + } + + // Clamp success score. + if successScore < 0 { + successScore = 0 + } + if successScore > 1 { + successScore = 1 + } + + // Complete via existing swarm service (validates assignee & status). + if err := s.swarm.CompleteTask(ctx, taskID, agentName); err != nil { + return nil, nil, err + } + + // Parse meta from the original task. + var meta AuctionTaskMeta + if len(task.Requirements) > 0 { + _ = json.Unmarshal(task.Requirements, &meta) + } + if meta.DifficultyWeight == 0 { + meta.DifficultyWeight = 1.0 + } + + // Find the winning bid to pick up estimated_tokens (if available). + var estimatedTokens int64 + for _, b := range bids { + if b.Status == channels.BidStatusAccepted { + var bm BidMeta + if len(b.Capabilities) > 0 { + _ = json.Unmarshal(b.Capabilities, &bm) + } + estimatedTokens = bm.EstimatedTokens + break + } + } + + // Default to a single "general" domain if none declared, so reputation + // is still recorded (FR-012 / SC-010). + domains := meta.Domains + if len(domains) == 0 { + domains = []string{"general"} + } + + var entries []*ReputationEntry + for _, d := range domains { + entry := &ReputationEntry{ + AgentName: agentName, + Domain: d, + TaskID: taskID, + EstimatedTokens: estimatedTokens, + ActualTokens: actualTokens, + SuccessScore: successScore, + DifficultyWeight: meta.DifficultyWeight, + } + if err := s.store.RecordEntry(ctx, entry); err != nil { + return nil, entries, fmt.Errorf("record reputation: %w", err) + } + entries = append(entries, entry) + } + + if s.tracer != nil { + s.tracer.Record(ctx, agentName, "marketplace.task_done", map[string]any{ + "task_id": taskID, + "actual_tokens": actualTokens, + "success_score": successScore, + "domains": domains, + }) + } + + task.Status = channels.TaskStatusCompleted + return task, entries, nil +} + +// QueryReputation returns the rolled-up reputation summary plus the most +// recent raw entries for (agent, domain). +func (s *Service) QueryReputation(ctx context.Context, agentName, domain string, limit int) (*ReputationSummary, []*ReputationEntry, error) { + if s.store == nil { + return nil, nil, fmt.Errorf("marketplace store not available") + } + sum, err := s.store.Summary(ctx, agentName, domain) + if err != nil { + return nil, nil, err + } + entries, err := s.store.ListEntries(ctx, agentName, domain, limit) + if err != nil { + return sum, nil, err + } + return sum, entries, nil +} diff --git a/internal/marketplace/store.go b/internal/marketplace/store.go new file mode 100644 index 0000000..d37137d --- /dev/null +++ b/internal/marketplace/store.go @@ -0,0 +1,196 @@ +// Package marketplace implements the agent marketplace MVP (spec 016): +// capability manifests (via wiki), auction channels (via existing swarm +// service), and the domain-scoped reputation ledger. +package marketplace + +import ( + "context" + "database/sql" + "fmt" + "time" +) + +// ReputationEntry is one row in the reputation ledger — a single +// (agent, domain) outcome from a completed auction task (FR-012). +type ReputationEntry struct { + ID int64 `json:"id"` + AgentName string `json:"agent_name"` + Domain string `json:"domain"` + TaskID int64 `json:"task_id,omitempty"` + EstimatedTokens int64 `json:"estimated_tokens"` + ActualTokens int64 `json:"actual_tokens"` + SuccessScore float64 `json:"success_score"` + DifficultyWeight float64 `json:"difficulty_weight"` + CompletedAt time.Time `json:"completed_at"` +} + +// ReputationSummary aggregates entries for one (agent, domain) pair. +type ReputationSummary struct { + AgentName string `json:"agent_name"` + Domain string `json:"domain"` + TasksCompleted int `json:"tasks_completed"` + AvgSuccessScore float64 `json:"avg_success_score"` + WeightedSuccessScore float64 `json:"weighted_success_score"` + AvgEstimatedTokens int64 `json:"avg_estimated_tokens"` + AvgActualTokens int64 `json:"avg_actual_tokens"` +} + +// Store persists reputation entries. +type Store struct { + db *sql.DB +} + +// NewStore creates a reputation store backed by the given SQLite handle. +func NewStore(db *sql.DB) *Store { + return &Store{db: db} +} + +// RecordEntry inserts one reputation row for (agent, domain) after a task completes. +func (s *Store) RecordEntry(ctx context.Context, e *ReputationEntry) error { + if e.AgentName == "" { + return fmt.Errorf("agent_name is required") + } + if e.Domain == "" { + return fmt.Errorf("domain is required") + } + if e.DifficultyWeight == 0 { + e.DifficultyWeight = 1.0 + } + + var taskID any + if e.TaskID > 0 { + taskID = e.TaskID + } + + res, err := s.db.ExecContext(ctx, + `INSERT INTO agent_reputation + (agent_name, domain, task_id, estimated_tokens, actual_tokens, + success_score, difficulty_weight, completed_at) + VALUES (?, ?, ?, ?, ?, ?, ?, CURRENT_TIMESTAMP)`, + e.AgentName, e.Domain, taskID, + e.EstimatedTokens, e.ActualTokens, + e.SuccessScore, e.DifficultyWeight, + ) + if err != nil { + return fmt.Errorf("insert reputation entry: %w", err) + } + id, err := res.LastInsertId() + if err != nil { + return fmt.Errorf("get reputation entry id: %w", err) + } + e.ID = id + return nil +} + +// ListEntries returns raw entries for (agent, domain) ordered by newest first. +// If domain is empty, all domains are returned. If agent is empty, the query is +// rejected — reputation is always scoped to an agent (FR-013). +func (s *Store) ListEntries(ctx context.Context, agent, domain string, limit int) ([]*ReputationEntry, error) { + if agent == "" { + return nil, fmt.Errorf("agent is required") + } + if limit <= 0 || limit > 500 { + limit = 100 + } + + var rows *sql.Rows + var err error + if domain == "" { + rows, err = s.db.QueryContext(ctx, + `SELECT id, agent_name, domain, COALESCE(task_id,0), + estimated_tokens, actual_tokens, success_score, + difficulty_weight, completed_at + FROM agent_reputation + WHERE agent_name = ? + ORDER BY completed_at DESC + LIMIT ?`, + agent, limit, + ) + } else { + rows, err = s.db.QueryContext(ctx, + `SELECT id, agent_name, domain, COALESCE(task_id,0), + estimated_tokens, actual_tokens, success_score, + difficulty_weight, completed_at + FROM agent_reputation + WHERE agent_name = ? AND domain = ? + ORDER BY completed_at DESC + LIMIT ?`, + agent, domain, limit, + ) + } + if err != nil { + return nil, fmt.Errorf("list reputation: %w", err) + } + defer rows.Close() + + var out []*ReputationEntry + for rows.Next() { + var e ReputationEntry + if err := rows.Scan(&e.ID, &e.AgentName, &e.Domain, &e.TaskID, + &e.EstimatedTokens, &e.ActualTokens, &e.SuccessScore, + &e.DifficultyWeight, &e.CompletedAt); err != nil { + return nil, fmt.Errorf("scan reputation row: %w", err) + } + out = append(out, &e) + } + if out == nil { + out = []*ReputationEntry{} + } + return out, rows.Err() +} + +// Summary computes aggregates for (agent, domain). Returns zero-values if +// there are no entries, and a count of 0. +func (s *Store) Summary(ctx context.Context, agent, domain string) (*ReputationSummary, error) { + if agent == "" { + return nil, fmt.Errorf("agent is required") + } + if domain == "" { + return nil, fmt.Errorf("domain is required") + } + + var ( + count int + sumSuccess sql.NullFloat64 + sumWeightedScore sql.NullFloat64 + sumWeights sql.NullFloat64 + sumEstimated sql.NullFloat64 + sumActual sql.NullFloat64 + ) + err := s.db.QueryRowContext(ctx, + `SELECT + COUNT(*), + SUM(success_score), + SUM(success_score * difficulty_weight), + SUM(difficulty_weight), + SUM(estimated_tokens), + SUM(actual_tokens) + FROM agent_reputation + WHERE agent_name = ? AND domain = ?`, + agent, domain, + ).Scan(&count, &sumSuccess, &sumWeightedScore, &sumWeights, &sumEstimated, &sumActual) + if err != nil { + return nil, fmt.Errorf("reputation summary: %w", err) + } + + sum := &ReputationSummary{ + AgentName: agent, + Domain: domain, + TasksCompleted: count, + } + if count > 0 { + if sumSuccess.Valid { + sum.AvgSuccessScore = sumSuccess.Float64 / float64(count) + } + if sumWeights.Valid && sumWeights.Float64 > 0 && sumWeightedScore.Valid { + sum.WeightedSuccessScore = sumWeightedScore.Float64 / sumWeights.Float64 + } + if sumEstimated.Valid { + sum.AvgEstimatedTokens = int64(sumEstimated.Float64 / float64(count)) + } + if sumActual.Valid { + sum.AvgActualTokens = int64(sumActual.Float64 / float64(count)) + } + } + return sum, nil +} diff --git a/internal/mcp/bridge.go b/internal/mcp/bridge.go index c633d5d..7efa2cd 100644 --- a/internal/mcp/bridge.go +++ b/internal/mcp/bridge.go @@ -15,6 +15,7 @@ import ( "github.com/synapbus/synapbus/internal/agentquery" "github.com/synapbus/synapbus/internal/attachments" "github.com/synapbus/synapbus/internal/channels" + "github.com/synapbus/synapbus/internal/marketplace" "github.com/synapbus/synapbus/internal/messaging" "github.com/synapbus/synapbus/internal/reactions" "github.com/synapbus/synapbus/internal/search" @@ -34,6 +35,7 @@ type ServiceBridge struct { reactionService *reactions.Service trustService *trust.Service wikiService *wiki.Service + marketplace *marketplace.Service queryExecutor *agentquery.Executor agentName string } @@ -156,6 +158,20 @@ func (b *ServiceBridge) Call(ctx context.Context, actionName string, args map[st case "send_message": return b.callSendMessage(ctx, args) + // --- Marketplace (spec 016) --- + case "post_auction": + return b.callPostAuction(ctx, args) + case "bid": + return b.callBid(ctx, args) + case "award": + return b.callAward(ctx, args) + case "mark_task_done": + return b.callMarkTaskDone(ctx, args) + case "read_skill_card": + return b.callReadSkillCard(ctx, args) + case "query_reputation": + return b.callQueryReputation(ctx, args) + default: return nil, fmt.Errorf("unknown action: %s", actionName) } diff --git a/internal/mcp/marketplace.go b/internal/mcp/marketplace.go new file mode 100644 index 0000000..5fec558 --- /dev/null +++ b/internal/mcp/marketplace.go @@ -0,0 +1,296 @@ +package mcp + +// Agent marketplace MCP actions (spec 016 MVP — US1, US2, US3). +// +// New actions exposed via the execute tool's call(actionName, args) interface: +// +// - post_auction — create an auction task with marketplace metadata +// - bid — place a bid with estimated_tokens / confidence +// - award — accept a winning bid and notify the winner +// - read_skill_card — read an agent's capability manifest (wiki article) +// - query_reputation — query the reputation ledger by (agent, domain) +// +// These actions follow the exact dispatch pattern used by bridge.go: +// each call* method validates args, calls the marketplace service, and +// returns a JSON-serialisable map. + +import ( + "context" + "encoding/json" + "fmt" + "strings" + "time" + + "github.com/synapbus/synapbus/internal/marketplace" +) + +// attachMarketplace installs the marketplace service on a ServiceBridge. +// Call it before dispatching any marketplace action. +func (b *ServiceBridge) attachMarketplace(mkt *marketplace.Service) { + b.marketplace = mkt +} + +func (b *ServiceBridge) callPostAuction(ctx context.Context, args map[string]any) (any, error) { + if b.marketplace == nil { + return nil, fmt.Errorf("marketplace service not available") + } + if b.channelService == nil { + return nil, fmt.Errorf("channel service not available") + } + + channelName := getString(args, "channel_name", "") + if channelName == "" { + return nil, fmt.Errorf("'channel_name' parameter is required") + } + title := getString(args, "title", "") + if title == "" { + return nil, fmt.Errorf("'title' parameter is required") + } + description := getString(args, "description", "") + acceptance := getString(args, "acceptance_criteria", "") + + meta := marketplace.AuctionTaskMeta{ + AcceptanceCriteria: acceptance, + MaxBudgetTokens: int64(getInt(args, "max_budget_tokens", 0)), + DifficultyWeight: getFloat(args, "difficulty_weight", 1.0), + } + + // domains: accept comma-separated string or []any + meta.Domains = parseStringList(args, "domains") + + deadlineStr := getString(args, "deadline", "") + var deadline *time.Time + if deadlineStr != "" { + t, err := time.Parse(time.RFC3339, deadlineStr) + if err != nil { + return nil, fmt.Errorf("deadline must be ISO 8601 format: %s", err) + } + deadline = &t + } + + ch, err := b.channelService.GetChannelByName(ctx, channelName) + if err != nil { + return nil, err + } + + task, err := b.marketplace.PostAuction(ctx, ch.ID, b.agentName, title, description, meta, deadline) + if err != nil { + return nil, err + } + + return map[string]any{ + "task_id": task.ID, + "channel_id": task.ChannelID, + "channel_name": channelName, + "title": task.Title, + "status": task.Status, + "posted_by": task.PostedBy, + "domains": meta.Domains, + "max_budget_tokens": meta.MaxBudgetTokens, + "deadline": task.Deadline, + "created_at": task.CreatedAt, + }, nil +} + +func (b *ServiceBridge) callBid(ctx context.Context, args map[string]any) (any, error) { + if b.marketplace == nil { + return nil, fmt.Errorf("marketplace service not available") + } + + taskID := int64(getInt(args, "task_id", 0)) + if taskID == 0 { + return nil, fmt.Errorf("'task_id' parameter is required") + } + + meta := marketplace.BidMeta{ + EstimatedTokens: int64(getInt(args, "estimated_tokens", 0)), + Confidence: getFloat(args, "confidence", 0), + Approach: getString(args, "approach", ""), + ManifestRevision: getInt(args, "manifest_revision", 0), + } + timeEstimate := getString(args, "time_estimate", "") + + bid, err := b.marketplace.Bid(ctx, taskID, b.agentName, meta, timeEstimate) + if err != nil { + return nil, err + } + return map[string]any{ + "bid_id": bid.ID, + "task_id": bid.TaskID, + "agent_name": bid.AgentName, + "status": bid.Status, + "estimated_tokens": meta.EstimatedTokens, + "confidence": meta.Confidence, + "approach": meta.Approach, + }, nil +} + +func (b *ServiceBridge) callAward(ctx context.Context, args map[string]any) (any, error) { + if b.marketplace == nil { + return nil, fmt.Errorf("marketplace service not available") + } + taskID := int64(getInt(args, "task_id", 0)) + if taskID == 0 { + return nil, fmt.Errorf("'task_id' parameter is required") + } + bidID := int64(getInt(args, "bid_id", 0)) + if bidID == 0 { + return nil, fmt.Errorf("'bid_id' parameter is required") + } + + winner, claimMsgID, err := b.marketplace.Award(ctx, taskID, bidID, b.agentName) + if err != nil { + return nil, err + } + return map[string]any{ + "task_id": taskID, + "bid_id": bidID, + "winner": winner, + "claim_message_id": claimMsgID, + "status": "awarded", + }, nil +} + +func (b *ServiceBridge) callMarkTaskDone(ctx context.Context, args map[string]any) (any, error) { + if b.marketplace == nil { + return nil, fmt.Errorf("marketplace service not available") + } + taskID := int64(getInt(args, "task_id", 0)) + if taskID == 0 { + return nil, fmt.Errorf("'task_id' parameter is required") + } + + actualTokens := int64(getInt(args, "actual_tokens", 0)) + successScore := getFloat(args, "success_score", 1.0) + + task, entries, err := b.marketplace.MarkTaskDone(ctx, taskID, b.agentName, actualTokens, successScore) + if err != nil { + return nil, err + } + return map[string]any{ + "task_id": taskID, + "status": task.Status, + "actual_tokens": actualTokens, + "success_score": successScore, + "reputation_entries": entries, + }, nil +} + +func (b *ServiceBridge) callReadSkillCard(ctx context.Context, args map[string]any) (any, error) { + if b.marketplace == nil { + return nil, fmt.Errorf("marketplace service not available") + } + + agentName := getString(args, "agent_name", "") + if agentName == "" { + agentName = b.agentName + } + + art, err := b.marketplace.ReadSkillCard(ctx, agentName) + if err != nil { + return nil, err + } + if art == nil { + return map[string]any{ + "agent_name": agentName, + "exists": false, + "slug": marketplace.SkillCardSlug(agentName), + }, nil + } + return map[string]any{ + "agent_name": agentName, + "exists": true, + "slug": art.Slug, + "title": art.Title, + "body": art.Body, + "revision": art.Revision, + "word_count": art.WordCount, + "updated_by": art.UpdatedBy, + "updated_at": art.UpdatedAt, + "outgoing_links": art.OutgoingLinks, + "backlinks": art.Backlinks, + }, nil +} + +func (b *ServiceBridge) callQueryReputation(ctx context.Context, args map[string]any) (any, error) { + if b.marketplace == nil { + return nil, fmt.Errorf("marketplace service not available") + } + agentName := getString(args, "agent_name", "") + if agentName == "" { + agentName = b.agentName + } + domain := getString(args, "domain", "") + if domain == "" { + return nil, fmt.Errorf("'domain' parameter is required (reputation is scoped by agent and domain)") + } + limit := getInt(args, "limit", 20) + + summary, entries, err := b.marketplace.QueryReputation(ctx, agentName, domain, limit) + if err != nil { + return nil, err + } + return map[string]any{ + "agent_name": agentName, + "domain": domain, + "summary": summary, + "recent_entries": entries, + "recent_count": len(entries), + }, nil +} + +// parseStringList reads args[key] as either a comma-separated string or a +// JSON/array-of-any and returns a trimmed list of non-empty strings. +func parseStringList(args map[string]any, key string) []string { + v, ok := args[key] + if !ok || v == nil { + return []string{} + } + var out []string + switch x := v.(type) { + case string: + // Accept either "a,b,c" or JSON-encoded ["a","b"] + s := strings.TrimSpace(x) + if s == "" { + return []string{} + } + if strings.HasPrefix(s, "[") { + var arr []string + if err := json.Unmarshal([]byte(s), &arr); err == nil { + for _, a := range arr { + a = strings.TrimSpace(a) + if a != "" { + out = append(out, a) + } + } + return out + } + } + for _, p := range strings.Split(s, ",") { + p = strings.TrimSpace(p) + if p != "" { + out = append(out, p) + } + } + case []any: + for _, item := range x { + if s, ok := item.(string); ok { + s = strings.TrimSpace(s) + if s != "" { + out = append(out, s) + } + } + } + case []string: + for _, s := range x { + s = strings.TrimSpace(s) + if s != "" { + out = append(out, s) + } + } + } + if out == nil { + out = []string{} + } + return out +} diff --git a/internal/mcp/marketplace_test.go b/internal/mcp/marketplace_test.go new file mode 100644 index 0000000..6823afa --- /dev/null +++ b/internal/mcp/marketplace_test.go @@ -0,0 +1,379 @@ +package mcp + +import ( + "context" + "testing" + + _ "modernc.org/sqlite" + + "github.com/synapbus/synapbus/internal/actions" + "github.com/synapbus/synapbus/internal/agents" + "github.com/synapbus/synapbus/internal/channels" + "github.com/synapbus/synapbus/internal/jsruntime" + "github.com/synapbus/synapbus/internal/marketplace" + "github.com/synapbus/synapbus/internal/messaging" + "github.com/synapbus/synapbus/internal/trace" + "github.com/synapbus/synapbus/internal/wiki" +) + +// newMarketplaceFixture builds a fully-wired bridge for marketplace tests. +// It seeds four agents (a poster + three bidders) and returns helpers for +// creating channels. +func newMarketplaceFixture(t *testing.T) (*HybridToolRegistrar, *ServiceBridge, *channels.Service, *marketplace.Service) { + t.Helper() + db := newTestDB(t) + + tracer := trace.NewTracer(db) + t.Cleanup(func() { tracer.Close() }) + + msgStore := messaging.NewSQLiteMessageStore(db) + msgService := messaging.NewMessagingService(msgStore, tracer) + + agentStore := agents.NewSQLiteAgentStore(db) + agentService := agents.NewAgentService(agentStore, tracer) + + channelStore := channels.NewSQLiteChannelStore(db) + channelService := channels.NewService(channelStore, msgService, tracer) + + taskStore := channels.NewSQLiteTaskStore(db) + swarmService := channels.NewSwarmService(taskStore, channelStore, tracer) + + wikiService := wiki.NewService(db) + + mktStore := marketplace.NewStore(db) + mkt := marketplace.NewService(mktStore, wikiService, swarmService, channelService, msgService, tracer) + + jsPool := jsruntime.NewPool(2) + t.Cleanup(func() { jsPool.Close() }) + + actionRegistry := actions.NewRegistry() + actionIndex := actions.NewIndex(actionRegistry.List()) + + registrar := NewHybridToolRegistrar( + msgService, + agentService, + channelService, + swarmService, + nil, // attachmentService + nil, // searchService + nil, // reactionService + nil, // trustService + wikiService, + jsPool, + actionRegistry, + actionIndex, + db, + ) + registrar.SetMarketplaceService(mkt) + + // Seed agents. + ctx := context.Background() + for _, name := range []string{"poster-agent", "bidder-alpha", "bidder-beta", "bidder-gamma"} { + if _, _, err := agentService.Register(ctx, name, name, "ai", nil, 1); err != nil { + t.Fatalf("seed agent %s: %v", name, err) + } + } + + bridge := NewServiceBridge( + msgService, + agentService, + channelService, + swarmService, + nil, nil, nil, nil, + wikiService, + "poster-agent", + ) + bridge.attachMarketplace(mkt) + return registrar, bridge, channelService, mkt +} + +// withAgent returns a clone of the bridge bound to a different agent name. +func (b *ServiceBridge) withAgent(name string) *ServiceBridge { + clone := *b + clone.agentName = name + return &clone +} + +func createAuctionChannel(t *testing.T, svc *channels.Service, name string, members ...string) int64 { + t.Helper() + ctx := context.Background() + ch, err := svc.CreateChannel(ctx, channels.CreateChannelRequest{ + Name: name, + Type: channels.TypeAuction, + CreatedBy: members[0], + }) + if err != nil { + t.Fatalf("create channel: %v", err) + } + for _, m := range members[1:] { + if err := svc.JoinChannel(ctx, ch.ID, m); err != nil { + t.Fatalf("join channel %s by %s: %v", name, m, err) + } + } + return ch.ID +} + +// -------------- US2: capability manifest via wiki -------------- + +func TestMarketplace_SkillCard_CreateReadUpdate(t *testing.T) { + _, bridge, _, _ := newMarketplaceFixture(t) + ctx := context.Background() + + // 1. read_skill_card before any article exists → exists=false. + res, err := bridge.Call(ctx, "read_skill_card", map[string]any{ + "agent_name": "poster-agent", + }) + if err != nil { + t.Fatalf("read_skill_card (empty): %v", err) + } + m := res.(map[string]any) + if m["exists"].(bool) != false { + t.Errorf("expected exists=false before publishing, got %v", m["exists"]) + } + if m["slug"].(string) != "agent-poster-agent" { + t.Errorf("unexpected slug %q", m["slug"]) + } + + // 2. publish a manifest using the existing wiki action (create_article). + _, err = bridge.Call(ctx, "create_article", map[string]any{ + "slug": "agent-poster-agent", + "title": "poster-agent skill card", + "body": "# Skills\n- data-analysis\n- python\n\nExample: [[mcp-gateway-competitors]]", + }) + if err != nil { + t.Fatalf("create_article: %v", err) + } + + // 3. read_skill_card now returns the article. + res, err = bridge.Call(ctx, "read_skill_card", map[string]any{ + "agent_name": "poster-agent", + }) + if err != nil { + t.Fatalf("read_skill_card: %v", err) + } + m = res.(map[string]any) + if m["exists"].(bool) != true { + t.Fatalf("expected exists=true after publish") + } + if rev, _ := m["revision"].(int); rev != 1 { + t.Errorf("expected revision=1, got %v", m["revision"]) + } + + // 4. update the article → revision increments (versioning for free via wiki). + _, err = bridge.Call(ctx, "update_article", map[string]any{ + "slug": "agent-poster-agent", + "body": "# Skills\n- data-analysis\n- python\n- go", + }) + if err != nil { + t.Fatalf("update_article: %v", err) + } + res, _ = bridge.Call(ctx, "read_skill_card", map[string]any{"agent_name": "poster-agent"}) + m = res.(map[string]any) + if rev, _ := m["revision"].(int); rev != 2 { + t.Errorf("expected revision=2 after update, got %v", m["revision"]) + } +} + +// -------------- US1: full auction lifecycle -------------- + +func TestMarketplace_FullAuctionLifecycle(t *testing.T) { + _, bridge, svc, _ := newMarketplaceFixture(t) + ctx := context.Background() + + // Create an auction channel with the poster and two bidders joined. + channelName := "task-marketplace-01" + createAuctionChannel(t, svc, channelName, "poster-agent", "bidder-alpha", "bidder-beta") + + // Poster posts an auction task with marketplace metadata. + res, err := bridge.Call(ctx, "post_auction", map[string]any{ + "channel_name": channelName, + "title": "Q4 revenue analysis", + "description": "Trend analysis on Q4 revenue with two charts", + "max_budget_tokens": 8000, + "domains": "data-analysis,python", + "difficulty_weight": 1.5, + }) + if err != nil { + t.Fatalf("post_auction: %v", err) + } + taskInfo := res.(map[string]any) + taskID := taskInfo["task_id"].(int64) + if taskID == 0 { + t.Fatal("expected non-zero task_id") + } + if domains, _ := taskInfo["domains"].([]string); len(domains) != 2 { + t.Errorf("expected 2 domains, got %v", taskInfo["domains"]) + } + + // Two bidders submit bids. + alpha := bridge.withAgent("bidder-alpha") + beta := bridge.withAgent("bidder-beta") + + bidRes, err := alpha.Call(ctx, "bid", map[string]any{ + "task_id": taskID, + "estimated_tokens": 4200, + "confidence": 0.9, + "approach": "pandas + matplotlib", + "manifest_revision": 1, + }) + if err != nil { + t.Fatalf("alpha bid: %v", err) + } + alphaBidID := bidRes.(map[string]any)["bid_id"].(int64) + + bidRes, err = beta.Call(ctx, "bid", map[string]any{ + "task_id": taskID, + "estimated_tokens": 6000, + "confidence": 0.7, + "approach": "R + ggplot", + }) + if err != nil { + t.Fatalf("beta bid: %v", err) + } + _ = bidRes + + // Self-bidding should be rejected (FR-011). + if _, err := bridge.Call(ctx, "bid", map[string]any{ + "task_id": taskID, + "estimated_tokens": 500, + }); err == nil { + t.Error("expected self-bid to be rejected") + } + + // Poster awards alpha's bid. + awardRes, err := bridge.Call(ctx, "award", map[string]any{ + "task_id": taskID, + "bid_id": alphaBidID, + }) + if err != nil { + t.Fatalf("award: %v", err) + } + am := awardRes.(map[string]any) + if am["winner"].(string) != "bidder-alpha" { + t.Errorf("winner = %v, want bidder-alpha", am["winner"]) + } + if am["claim_message_id"].(int64) == 0 { + t.Error("expected non-zero claim_message_id (DM to winner)") + } + + // Winner marks the task done with actual_tokens. + doneRes, err := alpha.Call(ctx, "mark_task_done", map[string]any{ + "task_id": taskID, + "actual_tokens": 4500, + "success_score": 1.0, + }) + if err != nil { + t.Fatalf("mark_task_done: %v", err) + } + dm := doneRes.(map[string]any) + if dm["status"].(string) != channels.TaskStatusCompleted { + t.Errorf("status = %v, want completed", dm["status"]) + } + + // -------------- US3: reputation ledger -------------- + + // There should be 2 reputation entries — one per domain. + entries := dm["reputation_entries"].([]*marketplace.ReputationEntry) + if len(entries) != 2 { + t.Errorf("expected 2 reputation entries (one per domain), got %d", len(entries)) + } + foundDataAnalysis := false + for _, e := range entries { + if e.Domain == "data-analysis" { + foundDataAnalysis = true + if e.ActualTokens != 4500 { + t.Errorf("data-analysis entry actual_tokens = %d, want 4500", e.ActualTokens) + } + if e.EstimatedTokens != 4200 { + t.Errorf("data-analysis entry estimated_tokens = %d, want 4200", e.EstimatedTokens) + } + if e.DifficultyWeight != 1.5 { + t.Errorf("data-analysis entry difficulty_weight = %f, want 1.5", e.DifficultyWeight) + } + } + } + if !foundDataAnalysis { + t.Error("expected a data-analysis reputation entry") + } + + // query_reputation returns the summary. + qrRes, err := bridge.Call(ctx, "query_reputation", map[string]any{ + "agent_name": "bidder-alpha", + "domain": "data-analysis", + }) + if err != nil { + t.Fatalf("query_reputation: %v", err) + } + qm := qrRes.(map[string]any) + summary := qm["summary"].(*marketplace.ReputationSummary) + if summary.TasksCompleted != 1 { + t.Errorf("tasks_completed = %d, want 1", summary.TasksCompleted) + } + if summary.AvgActualTokens != 4500 { + t.Errorf("avg_actual_tokens = %d, want 4500", summary.AvgActualTokens) + } + if summary.AvgSuccessScore != 1.0 { + t.Errorf("avg_success_score = %f, want 1.0", summary.AvgSuccessScore) + } + + // query_reputation for a domain with no entries returns zero summary. + qrRes, err = bridge.Call(ctx, "query_reputation", map[string]any{ + "agent_name": "bidder-alpha", + "domain": "unknown-domain", + }) + if err != nil { + t.Fatalf("query_reputation (empty): %v", err) + } + qm = qrRes.(map[string]any) + summary = qm["summary"].(*marketplace.ReputationSummary) + if summary.TasksCompleted != 0 { + t.Errorf("empty-domain tasks_completed = %d, want 0", summary.TasksCompleted) + } + + // Missing domain is rejected. + if _, err := bridge.Call(ctx, "query_reputation", map[string]any{ + "agent_name": "bidder-alpha", + }); err == nil { + t.Error("expected error when domain is missing") + } +} + +// Ensure post_auction requires an auction-type channel (reuses swarm guard). +func TestMarketplace_PostAuction_RejectsNonAuctionChannel(t *testing.T) { + _, bridge, svc, _ := newMarketplaceFixture(t) + ctx := context.Background() + + ch, err := svc.CreateChannel(ctx, channels.CreateChannelRequest{ + Name: "not-auction", + Type: channels.TypeStandard, + CreatedBy: "poster-agent", + }) + if err != nil { + t.Fatalf("create channel: %v", err) + } + _ = ch + + _, err = bridge.Call(ctx, "post_auction", map[string]any{ + "channel_name": "not-auction", + "title": "nope", + "domains": "x", + }) + if err == nil { + t.Error("expected error when posting to a non-auction channel") + } +} + +// Quick smoke test to make sure SkillCardSlug normalises underscores/case. +func TestSkillCardSlug_Normalises(t *testing.T) { + cases := map[string]string{ + "alice": "agent-alice", + "Bob": "agent-bob", + "research_alpha": "agent-research-alpha", + " spaced ": "agent-spaced", + } + for in, want := range cases { + if got := marketplace.SkillCardSlug(in); got != want { + t.Errorf("SkillCardSlug(%q) = %q, want %q", in, got, want) + } + } +} diff --git a/internal/mcp/server.go b/internal/mcp/server.go index a59acdd..25e788f 100644 --- a/internal/mcp/server.go +++ b/internal/mcp/server.go @@ -18,6 +18,7 @@ import ( "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" @@ -212,6 +213,14 @@ func (s *MCPServer) SetQueryExecutor(exec *agentquery.Executor) { } } +// 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 diff --git a/internal/mcp/tools_hybrid.go b/internal/mcp/tools_hybrid.go index 9b838ac..7924d8e 100644 --- a/internal/mcp/tools_hybrid.go +++ b/internal/mcp/tools_hybrid.go @@ -18,6 +18,7 @@ import ( "github.com/synapbus/synapbus/internal/channels" "github.com/synapbus/synapbus/internal/jsruntime" "github.com/synapbus/synapbus/internal/agentquery" + "github.com/synapbus/synapbus/internal/marketplace" "github.com/synapbus/synapbus/internal/messaging" "github.com/synapbus/synapbus/internal/reactions" "github.com/synapbus/synapbus/internal/search" @@ -36,6 +37,7 @@ type HybridToolRegistrar struct { reactionService *reactions.Service trustService *trust.Service wikiService *wiki.Service + marketplaceSvc *marketplace.Service jsPool *jsruntime.Pool actionRegistry *actions.Registry actionIndex *actions.Index @@ -44,6 +46,13 @@ type HybridToolRegistrar struct { logger *slog.Logger } +// SetMarketplaceService attaches the marketplace service for the 5 new +// spec-016 actions (post_auction, bid, award, read_skill_card, query_reputation). +// Call this after NewHybridToolRegistrar. +func (h *HybridToolRegistrar) SetMarketplaceService(m *marketplace.Service) { + h.marketplaceSvc = m +} + // SetQueryExecutor sets the SQL query executor for all agent bridges. func (h *HybridToolRegistrar) SetQueryExecutor(exec *agentquery.Executor) { h.queryExecutor = exec @@ -510,6 +519,9 @@ func (h *HybridToolRegistrar) handleExecute(ctx context.Context, req mcplib.Call if h.queryExecutor != nil { bridge.SetQueryExecutor(h.queryExecutor) } + if h.marketplaceSvc != nil { + bridge.attachMarketplace(h.marketplaceSvc) + } result, err := h.jsPool.Execute(ctx, code, bridge, jsruntime.ExecuteOptions{ Timeout: timeout, diff --git a/internal/reactions/model.go b/internal/reactions/model.go index f8ec348..fc715f3 100644 --- a/internal/reactions/model.go +++ b/internal/reactions/model.go @@ -14,6 +14,9 @@ const ( ReactionInProgress = "in_progress" ReactionDone = "done" ReactionPublished = "published" + // ReactionAwarded is used on auction-channel bid messages to designate + // the winning bid (spec 016, FR-008). Added in feature 016. + ReactionAwarded = "awarded" ) // Workflow states (derived from reactions). @@ -24,6 +27,7 @@ const ( StateRejected = "rejected" StateDone = "done" StatePublished = "published" + StateAwarded = "awarded" ) // reactionPriority maps reaction types to their priority for state derivation. @@ -32,8 +36,9 @@ var reactionPriority = map[string]int{ ReactionApprove: 2, ReactionInProgress: 3, ReactionReject: 4, - ReactionDone: 5, - ReactionPublished: 6, + ReactionAwarded: 5, + ReactionDone: 6, + ReactionPublished: 7, } // reactionToState maps reaction types to workflow states. @@ -41,6 +46,7 @@ var reactionToState = map[string]string{ ReactionApprove: StateApproved, ReactionReject: StateRejected, ReactionInProgress: StateInProgress, + ReactionAwarded: StateAwarded, ReactionDone: StateDone, ReactionPublished: StatePublished, } @@ -57,7 +63,7 @@ const MaxReactionsPerMessage = 100 // Sentinel errors. var ( - ErrInvalidReaction = errors.New("invalid reaction type: must be one of approve, reject, in_progress, done, published") + ErrInvalidReaction = errors.New("invalid reaction type: must be one of approve, reject, in_progress, awarded, done, published") ErrReactionLimit = errors.New("maximum reactions per message (100) reached") ErrNotMember = errors.New("only channel members can react to messages") ) @@ -77,6 +83,7 @@ var ValidReactions = map[string]bool{ ReactionApprove: true, ReactionReject: true, ReactionInProgress: true, + ReactionAwarded: true, ReactionDone: true, ReactionPublished: true, } diff --git a/internal/storage/schema/018_agent_marketplace.sql b/internal/storage/schema/018_agent_marketplace.sql new file mode 100644 index 0000000..e938fbf --- /dev/null +++ b/internal/storage/schema/018_agent_marketplace.sql @@ -0,0 +1,51 @@ +-- 018: Agent marketplace — reputation ledger + 'awarded' reaction (spec 016) +-- +-- Records one row per completed auction task per (agent, domain). A task +-- with multiple domain tags produces multiple rows, so reputation is queryable +-- as a vector across domains (FR-013). + +-- Widen the reactions CHECK constraint to allow the new 'awarded' reaction +-- used on auction-channel bid messages (FR-008). SQLite cannot ALTER a CHECK +-- constraint in place, so we rebuild the table. +CREATE TABLE IF NOT EXISTS message_reactions_new ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + message_id INTEGER NOT NULL REFERENCES messages(id) ON DELETE CASCADE, + agent_name TEXT NOT NULL, + reaction TEXT NOT NULL CHECK(reaction IN ('approve', 'reject', 'in_progress', 'awarded', 'done', 'published')), + metadata TEXT NOT NULL DEFAULT '{}', + created_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, + UNIQUE(message_id, agent_name, reaction) +); + +INSERT INTO message_reactions_new (id, message_id, agent_name, reaction, metadata, created_at) +SELECT id, message_id, agent_name, reaction, metadata, created_at FROM message_reactions; + +DROP TABLE message_reactions; +ALTER TABLE message_reactions_new RENAME TO message_reactions; + +CREATE INDEX IF NOT EXISTS idx_reactions_message ON message_reactions(message_id); +CREATE INDEX IF NOT EXISTS idx_reactions_agent ON message_reactions(agent_name); +CREATE INDEX IF NOT EXISTS idx_reactions_type ON message_reactions(reaction); + +CREATE TABLE IF NOT EXISTS agent_reputation ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + agent_name TEXT NOT NULL, + domain TEXT NOT NULL, + task_id INTEGER, -- source task; nullable for seed/manual entries + estimated_tokens INTEGER NOT NULL DEFAULT 0, + actual_tokens INTEGER NOT NULL DEFAULT 0, + success_score REAL NOT NULL DEFAULT 1.0, -- 0.0 (failure) .. 1.0 (perfect) + difficulty_weight REAL NOT NULL DEFAULT 1.0, -- per-task difficulty multiplier + completed_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP +); + +CREATE INDEX IF NOT EXISTS idx_agent_reputation_agent_domain + ON agent_reputation(agent_name, domain); + +CREATE INDEX IF NOT EXISTS idx_agent_reputation_domain + ON agent_reputation(domain); + +CREATE INDEX IF NOT EXISTS idx_agent_reputation_task + ON agent_reputation(task_id); + +INSERT INTO schema_migrations (version) VALUES (18);