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>
This commit is contained in:
Algis Dumbris
2026-04-11 15:17:30 +03:00
co-authored by Claude Opus 4.6
parent 96db7c06a4
commit cda3365863
12 changed files with 1440 additions and 8 deletions
+7
View File
@@ -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())
+101 -2
View File
@@ -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-<name>'. Returns exists=false if the agent has not published a manifest yet. Use create_article / update_article on slug 'agent-<your-name>' 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"})`,
}},
},
}
}
+7 -3
View File
@@ -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 {
+356
View File
@@ -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-<name>". 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
}
+196
View File
@@ -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
}
+16
View File
@@ -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)
}
+296
View File
@@ -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
}
+379
View File
@@ -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)
}
}
}
+9
View File
@@ -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
+12
View File
@@ -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,
+10 -3
View File
@@ -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,
}
@@ -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);