Slice 3: activity feed and active project (SSE + MCP tools)
Add internal/activity (thread-safe active project + bounded event feed with subscriber fan-out). Service exposes active-project/activity methods and takes the feed. New MCP tools get_active_project/set_active_project/get_activity and HTTP endpoints GET/POST /api/active-project, /api/activity, and /events (SSE). New <activity-feed> component updates live via EventSource; <repo-list> sets the active project on selection. Verified end-to-end in the running app: user actions push live events and set the active project (actor=user), all readable by Claude over MCP. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
@@ -0,0 +1,154 @@
|
||||
// Package activity holds the app's coordination state: the single active project
|
||||
// (the repo/task currently in focus) and a bounded feed of what happened — user
|
||||
// AND Claude actions. Both are in-memory (mirrored to the logs, no datastore —
|
||||
// AGENT.md §1.3) and queryable so Claude can sync on any turn boundary; new
|
||||
// events also fan out to subscribers for the browser SSE stream (§8.2). This is
|
||||
// the foundation the graceful project handoff (§8.3) builds on.
|
||||
package activity
|
||||
|
||||
import (
|
||||
"log/slog"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
// Actor is who caused an event.
|
||||
type Actor string
|
||||
|
||||
const (
|
||||
ActorUser Actor = "user"
|
||||
ActorClaude Actor = "claude"
|
||||
ActorSystem Actor = "system"
|
||||
)
|
||||
|
||||
// Event is one entry in the activity feed.
|
||||
type Event struct {
|
||||
ID int64 `json:"id"`
|
||||
Time time.Time `json:"time"`
|
||||
Actor Actor `json:"actor"`
|
||||
Kind string `json:"kind"` // e.g. "active-project-changed"
|
||||
Repo string `json:"repo,omitempty"` // repo path, when relevant
|
||||
Detail string `json:"detail,omitempty"` // human-readable extra context
|
||||
}
|
||||
|
||||
// Feed is the concurrency-safe active-project + activity store with fan-out.
|
||||
type Feed struct {
|
||||
mu sync.RWMutex
|
||||
active string // active project path ("" = none)
|
||||
events []Event
|
||||
maxEvents int
|
||||
nextID int64
|
||||
subs map[chan Event]struct{}
|
||||
log *slog.Logger
|
||||
}
|
||||
|
||||
// New builds a Feed keeping at most maxEvents recent events.
|
||||
func New(log *slog.Logger, maxEvents int) *Feed {
|
||||
if maxEvents <= 0 {
|
||||
maxEvents = 200
|
||||
}
|
||||
return &Feed{
|
||||
maxEvents: maxEvents,
|
||||
nextID: 1,
|
||||
subs: make(map[chan Event]struct{}),
|
||||
log: log,
|
||||
}
|
||||
}
|
||||
|
||||
// ActiveProject returns the current active project path ("" if none).
|
||||
func (f *Feed) ActiveProject() string {
|
||||
f.mu.RLock()
|
||||
defer f.mu.RUnlock()
|
||||
return f.active
|
||||
}
|
||||
|
||||
// SetActiveProject sets the active project and records an event. It is a no-op
|
||||
// (changed=false, zero Event) when path already matches, so repeated sets don't
|
||||
// spam the feed.
|
||||
func (f *Feed) SetActiveProject(actor Actor, path string) (Event, bool) {
|
||||
f.mu.Lock()
|
||||
if f.active == path {
|
||||
f.mu.Unlock()
|
||||
return Event{}, false
|
||||
}
|
||||
f.active = path
|
||||
ev, subs := f.appendLocked(actor, "active-project-changed", path, "")
|
||||
f.mu.Unlock()
|
||||
|
||||
publish(subs, ev)
|
||||
return ev, true
|
||||
}
|
||||
|
||||
// Record adds an arbitrary event to the feed.
|
||||
func (f *Feed) Record(actor Actor, kind, repo, detail string) Event {
|
||||
f.mu.Lock()
|
||||
ev, subs := f.appendLocked(actor, kind, repo, detail)
|
||||
f.mu.Unlock()
|
||||
|
||||
publish(subs, ev)
|
||||
return ev
|
||||
}
|
||||
|
||||
// Events returns up to limit of the most recent events, oldest first. limit<=0
|
||||
// returns all retained events.
|
||||
func (f *Feed) Events(limit int) []Event {
|
||||
f.mu.RLock()
|
||||
defer f.mu.RUnlock()
|
||||
if limit <= 0 || limit > len(f.events) {
|
||||
limit = len(f.events)
|
||||
}
|
||||
out := make([]Event, limit)
|
||||
copy(out, f.events[len(f.events)-limit:])
|
||||
return out
|
||||
}
|
||||
|
||||
// Subscribe returns a channel of future events and an unsubscribe func the
|
||||
// caller MUST invoke when done (e.g. via defer) to avoid leaking the channel.
|
||||
func (f *Feed) Subscribe() (<-chan Event, func()) {
|
||||
ch := make(chan Event, 16)
|
||||
f.mu.Lock()
|
||||
f.subs[ch] = struct{}{}
|
||||
f.mu.Unlock()
|
||||
|
||||
var once sync.Once
|
||||
unsub := func() {
|
||||
once.Do(func() {
|
||||
f.mu.Lock()
|
||||
delete(f.subs, ch)
|
||||
f.mu.Unlock()
|
||||
close(ch)
|
||||
})
|
||||
}
|
||||
return ch, unsub
|
||||
}
|
||||
|
||||
// appendLocked assigns id/time, appends (trimming to maxEvents), logs, and
|
||||
// returns the event plus a snapshot of subscriber channels to publish to after
|
||||
// the lock is released. Caller must hold f.mu.
|
||||
func (f *Feed) appendLocked(actor Actor, kind, repo, detail string) (Event, []chan Event) {
|
||||
ev := Event{ID: f.nextID, Time: time.Now(), Actor: actor, Kind: kind, Repo: repo, Detail: detail}
|
||||
f.nextID++
|
||||
f.events = append(f.events, ev)
|
||||
if len(f.events) > f.maxEvents {
|
||||
f.events = f.events[len(f.events)-f.maxEvents:]
|
||||
}
|
||||
if f.log != nil {
|
||||
f.log.Info("activity", "actor", actor, "kind", kind, "repo", repo, "detail", detail)
|
||||
}
|
||||
subs := make([]chan Event, 0, len(f.subs))
|
||||
for ch := range f.subs {
|
||||
subs = append(subs, ch)
|
||||
}
|
||||
return ev, subs
|
||||
}
|
||||
|
||||
// publish does a non-blocking send to each subscriber; a full channel (slow
|
||||
// consumer) drops the event rather than stalling the producer.
|
||||
func publish(subs []chan Event, ev Event) {
|
||||
for _, ch := range subs {
|
||||
select {
|
||||
case ch <- ev:
|
||||
default:
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -12,6 +12,7 @@ import (
|
||||
|
||||
mcpsdk "github.com/modelcontextprotocol/go-sdk/mcp"
|
||||
|
||||
"gitmanager/internal/activity"
|
||||
"gitmanager/internal/repos"
|
||||
"gitmanager/internal/service"
|
||||
)
|
||||
@@ -28,6 +29,21 @@ type listReposOutput struct {
|
||||
Repos []repos.State `json:"repos" jsonschema:"the discovered repositories"`
|
||||
}
|
||||
|
||||
// setActiveProjectInput is the argument schema for set_active_project.
|
||||
type setActiveProjectInput struct {
|
||||
Path string `json:"path" jsonschema:"absolute path of the repository to make active, exactly as returned by list_repos"`
|
||||
}
|
||||
|
||||
// activeProjectOutput reports the active project path (object, per the rule above).
|
||||
type activeProjectOutput struct {
|
||||
Path string `json:"path" jsonschema:"absolute path of the active project, empty when none is set"`
|
||||
}
|
||||
|
||||
// activityOutput wraps the activity feed (object, per the rule above).
|
||||
type activityOutput struct {
|
||||
Events []activity.Event `json:"events" jsonschema:"recent activity events, oldest first"`
|
||||
}
|
||||
|
||||
// NewServer builds the MCP server and registers the (currently read-only) tools.
|
||||
func NewServer(svc *service.Service, version string) *mcpsdk.Server {
|
||||
s := mcpsdk.NewServer(&mcpsdk.Implementation{
|
||||
@@ -57,6 +73,33 @@ func NewServer(svc *service.Service, version string) *mcpsdk.Server {
|
||||
return nil, detail, nil
|
||||
})
|
||||
|
||||
// get_active_project — the repo/task currently in focus (§8.2).
|
||||
mcpsdk.AddTool(s, &mcpsdk.Tool{
|
||||
Name: "get_active_project",
|
||||
Description: "Get the active project — the repository the user is currently focused on. Check this to stay in sync with the user; path is empty when none is set.",
|
||||
}, func(_ context.Context, _ *mcpsdk.CallToolRequest, _ struct{}) (*mcpsdk.CallToolResult, activeProjectOutput, error) {
|
||||
return nil, activeProjectOutput{Path: svc.ActiveProject()}, nil
|
||||
})
|
||||
|
||||
// set_active_project — Claude switches the focus to another repo.
|
||||
mcpsdk.AddTool(s, &mcpsdk.Tool{
|
||||
Name: "set_active_project",
|
||||
Description: "Set the active project to the given repository path (from list_repos). Use this when switching which repository you are working in so the app and user stay in sync.",
|
||||
}, func(_ context.Context, _ *mcpsdk.CallToolRequest, in setActiveProjectInput) (*mcpsdk.CallToolResult, activeProjectOutput, error) {
|
||||
if _, _, err := svc.SetActiveProject(activity.ActorClaude, in.Path); err != nil {
|
||||
return nil, activeProjectOutput{}, fmt.Errorf("%w — call list_repos for valid paths", err)
|
||||
}
|
||||
return nil, activeProjectOutput{Path: svc.ActiveProject()}, nil
|
||||
})
|
||||
|
||||
// get_activity — recent user + Claude actions, so Claude can catch up.
|
||||
mcpsdk.AddTool(s, &mcpsdk.Tool{
|
||||
Name: "get_activity",
|
||||
Description: "Get the recent activity feed (user and Claude actions, oldest first): repo selections, active-project changes, and more as features land. Use it to see what the user has done since you last looked.",
|
||||
}, func(_ context.Context, _ *mcpsdk.CallToolRequest, _ struct{}) (*mcpsdk.CallToolResult, activityOutput, error) {
|
||||
return nil, activityOutput{Events: svc.Activity(50)}, nil
|
||||
})
|
||||
|
||||
return s
|
||||
}
|
||||
|
||||
|
||||
@@ -13,6 +13,7 @@ import (
|
||||
|
||||
mcpsdk "github.com/modelcontextprotocol/go-sdk/mcp"
|
||||
|
||||
"gitmanager/internal/activity"
|
||||
"gitmanager/internal/git"
|
||||
"gitmanager/internal/repos"
|
||||
"gitmanager/internal/service"
|
||||
@@ -41,7 +42,7 @@ func TestMCPRoundTrip(t *testing.T) {
|
||||
scanner := repos.NewScanner(g, log, []string{root}, 3, nil, time.Minute, false)
|
||||
scanner.Refresh(context.Background())
|
||||
|
||||
svc := service.New(g, scanner.Index)
|
||||
svc := service.New(g, scanner.Index, activity.New(log, 200))
|
||||
srv := NewServer(svc, "test")
|
||||
|
||||
// Wire an in-memory client<->server session.
|
||||
@@ -100,6 +101,64 @@ func TestMCPRoundTrip(t *testing.T) {
|
||||
if !res.IsError {
|
||||
t.Fatalf("expected IsError for unknown repo, got success")
|
||||
}
|
||||
|
||||
// active project starts empty.
|
||||
res, err = cs.CallTool(ctx, &mcpsdk.CallToolParams{Name: "get_active_project"})
|
||||
if err != nil {
|
||||
t.Fatalf("get_active_project: %v", err)
|
||||
}
|
||||
var ap struct {
|
||||
Path string `json:"path"`
|
||||
}
|
||||
decodeResult(t, res, &ap)
|
||||
if ap.Path != "" {
|
||||
t.Fatalf("expected empty active project, got %q", ap.Path)
|
||||
}
|
||||
|
||||
// set_active_project to our repo, then read it back.
|
||||
res, err = cs.CallTool(ctx, &mcpsdk.CallToolParams{
|
||||
Name: "set_active_project",
|
||||
Arguments: map[string]any{"path": states[0].Path},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("set_active_project: %v", err)
|
||||
}
|
||||
if res.IsError {
|
||||
t.Fatalf("set_active_project returned tool error: %+v", res.Content)
|
||||
}
|
||||
res, _ = cs.CallTool(ctx, &mcpsdk.CallToolParams{Name: "get_active_project"})
|
||||
decodeResult(t, res, &ap)
|
||||
if ap.Path != states[0].Path {
|
||||
t.Fatalf("active project = %q, want %q", ap.Path, states[0].Path)
|
||||
}
|
||||
|
||||
// setting an unknown project is a tool error.
|
||||
res, err = cs.CallTool(ctx, &mcpsdk.CallToolParams{
|
||||
Name: "set_active_project",
|
||||
Arguments: map[string]any{"path": filepath.Join(root, "nope")},
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("set_active_project(bad) protocol error: %v", err)
|
||||
}
|
||||
if !res.IsError {
|
||||
t.Fatalf("expected IsError for unknown active project")
|
||||
}
|
||||
|
||||
// the activity feed should now contain the active-project-changed event.
|
||||
res, _ = cs.CallTool(ctx, &mcpsdk.CallToolParams{Name: "get_activity"})
|
||||
var act struct {
|
||||
Events []activity.Event `json:"events"`
|
||||
}
|
||||
decodeResult(t, res, &act)
|
||||
found := false
|
||||
for _, ev := range act.Events {
|
||||
if ev.Kind == "active-project-changed" && ev.Repo == states[0].Path && ev.Actor == activity.ActorClaude {
|
||||
found = true
|
||||
}
|
||||
}
|
||||
if !found {
|
||||
t.Fatalf("expected an active-project-changed event by claude, got %+v", act.Events)
|
||||
}
|
||||
}
|
||||
|
||||
// decodeResult unmarshals the JSON text content of a tool result into v.
|
||||
|
||||
@@ -7,8 +7,10 @@ package service
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"path/filepath"
|
||||
|
||||
"gitmanager/internal/activity"
|
||||
"gitmanager/internal/git"
|
||||
"gitmanager/internal/repos"
|
||||
)
|
||||
@@ -17,11 +19,13 @@ import (
|
||||
type Service struct {
|
||||
git *git.CLI
|
||||
index *repos.Index
|
||||
feed *activity.Feed
|
||||
}
|
||||
|
||||
// New builds a Service over the git boundary and the scanner's repo index.
|
||||
func New(g *git.CLI, index *repos.Index) *Service {
|
||||
return &Service{git: g, index: index}
|
||||
// New builds a Service over the git boundary, the scanner's repo index, and the
|
||||
// activity feed.
|
||||
func New(g *git.CLI, index *repos.Index, feed *activity.Feed) *Service {
|
||||
return &Service{git: g, index: index, feed: feed}
|
||||
}
|
||||
|
||||
// ListRepos returns a snapshot of every discovered repository.
|
||||
@@ -45,3 +49,40 @@ func (s *Service) RepoDetail(ctx context.Context, path string) (repos.Detail, bo
|
||||
}
|
||||
return repos.BuildDetail(ctx, s.git, base), true
|
||||
}
|
||||
|
||||
// --- Activity & active project (§8.2) --------------------------------------
|
||||
|
||||
// ActiveProject returns the current active project path ("" if none).
|
||||
func (s *Service) ActiveProject() string {
|
||||
return s.feed.ActiveProject()
|
||||
}
|
||||
|
||||
// SetActiveProject makes path the active project (or clears it when empty). It
|
||||
// rejects a path that is not an indexed repository — the active project must be
|
||||
// a real repo (§1.3). Returns the recorded event and whether it changed.
|
||||
func (s *Service) SetActiveProject(actor activity.Actor, path string) (activity.Event, bool, error) {
|
||||
if path != "" {
|
||||
path = filepath.Clean(path)
|
||||
if _, ok := s.index.Get(path); !ok {
|
||||
return activity.Event{}, false, fmt.Errorf("unknown repository %q", path)
|
||||
}
|
||||
}
|
||||
ev, changed := s.feed.SetActiveProject(actor, path)
|
||||
return ev, changed, nil
|
||||
}
|
||||
|
||||
// RecordActivity appends an arbitrary event to the feed.
|
||||
func (s *Service) RecordActivity(actor activity.Actor, kind, repo, detail string) activity.Event {
|
||||
return s.feed.Record(actor, kind, repo, detail)
|
||||
}
|
||||
|
||||
// Activity returns up to limit recent events, oldest first.
|
||||
func (s *Service) Activity(limit int) []activity.Event {
|
||||
return s.feed.Events(limit)
|
||||
}
|
||||
|
||||
// SubscribeActivity returns a channel of future events plus an unsubscribe func
|
||||
// the caller must invoke when done.
|
||||
func (s *Service) SubscribeActivity() (<-chan activity.Event, func()) {
|
||||
return s.feed.Subscribe()
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user