Files
Plexum-fabric-channel-plugin/internal/fabric/client.go
hzhang 7911cc6320 feat(presence): F-5b presence-sync — mirror sm.Machine into Fabric
Adds an internal/presence package that ticks every 30s (configurable
via presence_interval_seconds), reads each bound agent's sm.Machine
state through host.ReadAgentState, maps to Fabric's 6-status enum,
and PUTs diffs to /api/agents/:userId/presence on every guild the
agent belongs to.

Semantic mapping (the part flagged "needed" in the prior README):
  idle    → idle
  working → on_call
  busy    → busy
  offline → offline
exhausted/unknown reserved for backend-side fallbacks; we don't push.

Tick is mutex-guarded (avoids the upsert race the openclaw incident
called out in agent-presence.service.ts) and diff-gated so writes are
sparse. Token-cache invalidation on PUT failure handles guild JWT
rotation.

fabric.Client gains SetAgentPresence helper. README marks F-5b .
2026-06-01 08:40:17 +01:00

456 lines
15 KiB
Go

// Package fabric is a thin Go port of Fabric.OpenclawPlugin's
// fabric-client.ts — Center auth + Guild REST. v0.1 covers what
// F-1 needs (auth/login, refresh, me/guilds, postMessage, listChannels,
// listMessages, channelMembers); the canvas / commands / sub-discussion
// surfaces arrive in later phases as their tools land.
package fabric
import (
"bytes"
"context"
"encoding/json"
"fmt"
"io"
"net/http"
"net/url"
"time"
)
// Session is what /auth/agent/login returns.
type Session struct {
AccessToken string `json:"accessToken"`
RefreshToken string `json:"refreshToken"`
User SessionUser `json:"user"`
Guilds []GuildInfo `json:"guilds"`
GuildAccessTokens []GuildAccessToken `json:"guildAccessTokens"`
}
// SessionUser is the user metadata baked into the session.
type SessionUser struct {
ID string `json:"id"`
Email string `json:"email"`
Name string `json:"name"`
}
// GuildInfo describes one guild this user belongs to.
type GuildInfo struct {
NodeID string `json:"nodeId"`
Name string `json:"name"`
Endpoint string `json:"endpoint"`
Status string `json:"status"`
Purpose *string `json:"purpose,omitempty"`
}
// GuildAccessToken pairs a per-guild short-lived JWT with the guild node.
type GuildAccessToken struct {
GuildNodeID string `json:"guildNodeId"`
Token string `json:"token"`
}
// Client is a thin wrapper around net/http.Client.
type Client struct {
CenterAPIBase string // e.g. "http://localhost:7001/api"
HTTP *http.Client
}
// New constructs a Client with a sensible default http client (30s timeout).
func New(centerAPIBase string) *Client {
return &Client{
CenterAPIBase: centerAPIBase,
HTTP: &http.Client{Timeout: 30 * time.Second},
}
}
// ---- low-level helpers ----
func (c *Client) do(ctx context.Context, method, url string, auth string, body any, extraHeaders map[string]string) ([]byte, error) {
var reader io.Reader
if body != nil {
raw, err := json.Marshal(body)
if err != nil {
return nil, fmt.Errorf("fabric: marshal: %w", err)
}
reader = bytes.NewReader(raw)
}
req, err := http.NewRequestWithContext(ctx, method, url, reader)
if err != nil {
return nil, err
}
if body != nil {
req.Header.Set("content-type", "application/json")
}
if auth != "" {
req.Header.Set("authorization", "Bearer "+auth)
}
for k, v := range extraHeaders {
req.Header.Set(k, v)
}
resp, err := c.HTTP.Do(req)
if err != nil {
return nil, fmt.Errorf("fabric: %s %s: %w", method, url, err)
}
defer resp.Body.Close()
raw, _ := io.ReadAll(resp.Body)
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
return nil, fmt.Errorf("fabric: %s %s -> %d: %s", method, url, resp.StatusCode, string(raw))
}
return raw, nil
}
func (c *Client) postJSON(ctx context.Context, url string, body any, auth string, out any) error {
raw, err := c.do(ctx, http.MethodPost, url, auth, body, nil)
if err != nil {
return err
}
if out == nil || len(raw) == 0 {
return nil
}
return json.Unmarshal(raw, out)
}
func (c *Client) getJSON(ctx context.Context, url, auth string, out any) error {
raw, err := c.do(ctx, http.MethodGet, url, auth, nil, nil)
if err != nil {
return err
}
if out == nil || len(raw) == 0 {
return nil
}
return json.Unmarshal(raw, out)
}
// ---- Center: auth ----
// AgentLogin exchanges an API key for a fresh session + guild tokens.
func (c *Client) AgentLogin(ctx context.Context, apiKey string) (*Session, error) {
var s Session
if err := c.postJSON(ctx, c.CenterAPIBase+"/auth/agent/login",
map[string]string{"apiKey": apiKey}, "", &s); err != nil {
return nil, err
}
return &s, nil
}
// Refresh trades a refresh token for a fresh access token (guild
// tokens are re-fetched separately via MeGuilds).
func (c *Client) Refresh(ctx context.Context, refreshToken string) (*RefreshResponse, error) {
var out RefreshResponse
if err := c.postJSON(ctx, c.CenterAPIBase+"/auth/refresh",
map[string]string{"refreshToken": refreshToken}, "", &out); err != nil {
return nil, err
}
return &out, nil
}
// RefreshResponse is the shape returned by /auth/refresh.
type RefreshResponse struct {
AccessToken string `json:"accessToken"`
RefreshToken string `json:"refreshToken"`
}
// MeGuilds returns the calling user's guild list + fresh per-guild tokens.
func (c *Client) MeGuilds(ctx context.Context, accessToken string) (*MeGuildsResponse, error) {
var out MeGuildsResponse
if err := c.getJSON(ctx, c.CenterAPIBase+"/auth/me/guilds", accessToken, &out); err != nil {
return nil, err
}
return &out, nil
}
// MeGuildsResponse subset of Session.
type MeGuildsResponse struct {
Guilds []GuildInfo `json:"guilds"`
GuildAccessTokens []GuildAccessToken `json:"guildAccessTokens"`
}
// ---- Guild: messaging ----
// PostMessage posts plain content to a channel as authorUserID.
func (c *Client) PostMessage(ctx context.Context, guildEndpoint, guildToken, channelID, content, authorUserID string) error {
_, err := c.do(ctx, http.MethodPost,
guildEndpoint+"/api/channels/"+url.PathEscape(channelID)+"/messages",
guildToken,
map[string]string{"content": content, "authorUserId": authorUserID},
nil)
return err
}
// PostSystemMessage posts a system-kind message. Guild routes
// kind=sys differently (not subject to turn engine, doesn't wake
// agents); useful for narrator-style narration from a plugin tool.
func (c *Client) PostSystemMessage(ctx context.Context, guildEndpoint, guildToken, channelID, content, authorUserID string) error {
_, err := c.do(ctx, http.MethodPost,
guildEndpoint+"/api/channels/"+url.PathEscape(channelID)+"/messages",
guildToken,
map[string]any{"content": content, "authorUserId": authorUserID, "kind": "sys"},
nil)
return err
}
// CreateChannelOpts is the payload to POST /api/channels.
type CreateChannelOpts struct {
GuildID string `json:"guildId"`
Name string `json:"name"`
XType string `json:"xType"`
IsPublic bool `json:"isPublic,omitempty"`
MemberUserIDs []string `json:"memberUserIds,omitempty"`
OnDuty string `json:"onDuty,omitempty"`
Listeners []string `json:"listeners,omitempty"`
Purpose string `json:"purpose,omitempty"`
}
// CreateChannel creates a new channel in a guild. Returns the new channel id.
func (c *Client) CreateChannel(ctx context.Context, guildEndpoint, guildToken string, opts CreateChannelOpts) (string, error) {
var out struct {
ID string `json:"id"`
}
if err := c.postJSON(ctx, guildEndpoint+"/api/channels", opts, guildToken, &out); err != nil {
return "", err
}
return out.ID, nil
}
// CloseChannel closes a channel (one-way for most x_types).
func (c *Client) CloseChannel(ctx context.Context, guildEndpoint, guildToken, channelID string) error {
_, err := c.do(ctx, http.MethodPost,
guildEndpoint+"/api/channels/"+url.PathEscape(channelID)+"/close",
guildToken, map[string]any{}, nil)
return err
}
// JoinChannel adds the calling user to a public channel.
func (c *Client) JoinChannel(ctx context.Context, guildEndpoint, guildToken, channelID string) error {
_, err := c.do(ctx, http.MethodPost,
guildEndpoint+"/api/channels/"+url.PathEscape(channelID)+"/join",
guildToken, map[string]any{}, nil)
return err
}
// LeaveChannel removes the calling user from a channel.
func (c *Client) LeaveChannel(ctx context.Context, guildEndpoint, guildToken, channelID string) error {
_, err := c.do(ctx, http.MethodPost,
guildEndpoint+"/api/channels/"+url.PathEscape(channelID)+"/leave",
guildToken, map[string]any{}, nil)
return err
}
// SetChannelPurpose PATCHes the channel's purpose (free-form text).
func (c *Client) SetChannelPurpose(ctx context.Context, guildEndpoint, guildToken, channelID, purpose string) error {
_, err := c.do(ctx, http.MethodPatch,
guildEndpoint+"/api/channels/"+url.PathEscape(channelID),
guildToken, map[string]string{"purpose": purpose}, nil)
return err
}
// ---- channel canvas ----
// CanvasFormat is "md" | "html" | "text".
type CanvasFormat string
// Canvas is the shape returned by GET /api/channels/<id>/canvas.
type Canvas struct {
ChannelID string `json:"channelId"`
SharerUserID string `json:"sharerUserId"`
Title string `json:"title"`
Format CanvasFormat `json:"format"`
Source string `json:"source"`
UpdatedAt string `json:"updatedAt,omitempty"`
}
// GetCanvas returns the canvas or nil if no canvas is set.
func (c *Client) GetCanvas(ctx context.Context, guildEndpoint, guildToken, channelID string) (*Canvas, error) {
raw, err := c.do(ctx, http.MethodGet,
guildEndpoint+"/api/channels/"+url.PathEscape(channelID)+"/canvas",
guildToken, nil, nil)
if err != nil {
return nil, err
}
if len(raw) == 0 {
return nil, nil
}
var out Canvas
if err := json.Unmarshal(raw, &out); err != nil {
return nil, err
}
return &out, nil
}
// ShareCanvas replaces (or initializes) the channel canvas. Caller
// becomes the sharer.
func (c *Client) ShareCanvas(ctx context.Context, guildEndpoint, guildToken, channelID, title string, format CanvasFormat, source string) (*Canvas, error) {
body := map[string]any{"title": title, "format": format, "source": source}
raw, err := c.do(ctx, http.MethodPut,
guildEndpoint+"/api/channels/"+url.PathEscape(channelID)+"/canvas",
guildToken, body, nil)
if err != nil {
return nil, err
}
var out Canvas
if err := json.Unmarshal(raw, &out); err != nil {
return nil, err
}
return &out, nil
}
// UpdateCanvas updates fields in place (original sharer only; else 403).
func (c *Client) UpdateCanvas(ctx context.Context, guildEndpoint, guildToken, channelID string, title, source string, format CanvasFormat) (*Canvas, error) {
body := map[string]any{}
if title != "" {
body["title"] = title
}
if source != "" {
body["source"] = source
}
if format != "" {
body["format"] = format
}
raw, err := c.do(ctx, http.MethodPatch,
guildEndpoint+"/api/channels/"+url.PathEscape(channelID)+"/canvas",
guildToken, body, nil)
if err != nil {
return nil, err
}
var out Canvas
if err := json.Unmarshal(raw, &out); err != nil {
return nil, err
}
return &out, nil
}
// RemoveCanvas closes the canvas.
func (c *Client) RemoveCanvas(ctx context.Context, guildEndpoint, guildToken, channelID string) error {
_, err := c.do(ctx, http.MethodDelete,
guildEndpoint+"/api/channels/"+url.PathEscape(channelID)+"/canvas",
guildToken, nil, nil)
return err
}
// SetAgentPresence pushes one agent's presence to the guild. Body
// shape mirrors PUT /api/agents/:userId/presence:
//
// { "status": "idle"|"on_call"|"busy"|"exhausted"|"offline"|"unknown",
// "source": "<debug-tag>" }
//
// userID is the agent's Fabric Center user UUID (not the Plexum agent
// id). The guild-side ApiKeyGuard accepts the per-guild access JWT —
// caller must already hold a fresh one.
func (c *Client) SetAgentPresence(ctx context.Context, guildEndpoint, guildToken, userID, status, source string) error {
body := map[string]any{"status": status, "source": source}
_, err := c.do(ctx, http.MethodPut,
guildEndpoint+"/api/agents/"+url.PathEscape(userID)+"/presence",
guildToken, body, nil)
return err
}
// SyncCommands PUTs the agent's slash-command catalog onto the guild
// (idempotent full replace). Needs the guild's commands-sync key, which
// the operator sources from the guild config.
func (c *Client) SyncCommands(ctx context.Context, guildEndpoint, guildToken string, commands []any, syncKey string) error {
headers := map[string]string{}
if syncKey != "" {
headers["x-commands-sync-key"] = syncKey
}
_, err := c.do(ctx, http.MethodPut,
guildEndpoint+"/api/commands", guildToken,
map[string]any{"commands": commands}, headers)
return err
}
// ChannelMembers lists members of a channel.
type ChannelMember struct {
UserID string `json:"userId"`
Bypass bool `json:"bypass,omitempty"`
}
func (c *Client) ChannelMembers(ctx context.Context, guildEndpoint, guildToken, channelID string) ([]ChannelMember, error) {
var out []ChannelMember
if err := c.getJSON(ctx,
guildEndpoint+"/api/channels/"+url.PathEscape(channelID)+"/members",
guildToken, &out); err != nil {
return nil, err
}
return out, nil
}
// ---- Guild: channel discovery + history ----
// Channel is the wire shape returned by /api/channels list/get.
type Channel struct {
ID string `json:"id"`
GuildID string `json:"guildId"`
Name string `json:"name"`
XType string `json:"xType"`
Kind string `json:"kind"`
IsPublic bool `json:"isPublic"`
Closed bool `json:"closed"`
LastSeq int `json:"lastSeq"`
CreatedAt string `json:"createdAt"`
Purpose *string `json:"purpose,omitempty"`
}
// ListChannels lists all channels in a guild visible to the calling user.
func (c *Client) ListChannels(ctx context.Context, guildEndpoint, guildToken, guildNodeID string) ([]Channel, error) {
var out []Channel
u := guildEndpoint + "/api/channels?guildId=" + url.QueryEscape(guildNodeID)
if err := c.getJSON(ctx, u, guildToken, &out); err != nil {
return nil, err
}
return out, nil
}
// Message is the wire shape of one message in history.
type Message struct {
MessageID string `json:"messageId"`
Seq int `json:"seq"`
Content string `json:"content"`
AuthorUserID string `json:"authorUserId"`
CreatedAt string `json:"createdAt"`
EditedAt *string `json:"editedAt"`
DeletedAt *string `json:"deletedAt"`
IsDeleted bool `json:"isDeleted"`
}
// MessagePage wraps a window of messages + pagination metadata.
type MessagePage struct {
Items []Message `json:"items"`
Page struct {
SeqFrom int `json:"seqFrom"`
SeqTo int `json:"seqTo"`
Limit int `json:"limit"`
Returned int `json:"returned"`
HasMore bool `json:"hasMore"`
NextExpectedSeq int `json:"nextExpectedSeq"`
HighestCommittedSeq int `json:"highestCommittedSeq"`
} `json:"page"`
}
// ListMessages fetches a window of messages by seq.
func (c *Client) ListMessages(ctx context.Context, guildEndpoint, guildToken, channelID string, opts ListMessagesOpts) (*MessagePage, error) {
qs := url.Values{}
if opts.SeqFrom > 0 {
qs.Set("seq_from", fmt.Sprint(opts.SeqFrom))
}
if opts.SeqTo > 0 {
qs.Set("seq_to", fmt.Sprint(opts.SeqTo))
}
if opts.Limit > 0 {
qs.Set("limit", fmt.Sprint(opts.Limit))
}
u := guildEndpoint + "/api/channels/" + url.PathEscape(channelID) + "/messages"
if encoded := qs.Encode(); encoded != "" {
u += "?" + encoded
}
var out MessagePage
if err := c.getJSON(ctx, u, guildToken, &out); err != nil {
return nil, err
}
return &out, nil
}
// ListMessagesOpts is the optional paging window.
type ListMessagesOpts struct {
SeqFrom int
SeqTo int
Limit int
}