Files
multica/server/internal/daemonws/notifier.go
Multica Eve c25a82eee0 perf(agents): fast model discovery on runtime switch (MUL-5444) (#6098)
* perf(agents): make runtime model discovery fast on runtime switch (MUL-5444)

Switching runtime in the agent creation form left the model picker
spinning for ~8-20s. Two costs stacked up:

- the list-models request sat in the store until the daemon's next
  scheduled heartbeat (0-15s, avg 7.5s of pure dead wait), and
- the daemon then enumerated the catalog locally (static for claude,
  but a CLI/ACP round trip up to ~15s for everyone else).

Both are addressed with the two standard techniques for a slow,
low-frequency, read-only operation: push instead of poll, and
stale-while-revalidate.

Push (removes the heartbeat wait):
- new additive `daemon:pending_work` hint, runtime-scoped, delivered
  through the existing daemon WS hub and the Redis relay so the API node
  holding the socket does the delivery.
- the daemon answers a hint with ONE immediate heartbeat and dispatches
  what it claimed. The hint deliberately carries no work, so nothing has
  to be un-claimed when delivery fails and a duplicate hint cannot
  duplicate work - PopPending stays the atomic claim.
- per-runtime coalescing plus a 1s floor keeps a caller-triggered hint
  from becoming a heartbeat amplifier.

Cache (removes the discovery wait on repeat opens):
- server-side per-runtime catalog cache (in-memory single-node, Redis
  multi-node) written on every successful report.
- a snapshot younger than 15min answers the POST immediately as an
  already-completed request; older than 60s it also enqueues a
  background refresh that only warms the cache.
- only supported, non-empty catalogs are cached; a completed-but-empty
  report invalidates instead, while a failed report keeps serving the
  last known good list.

Frontend: staleTime 60s -> 5min and gcTime 30min, so a runtime revisited
in the same session renders from cache and revalidates in the background
instead of showing the spinner again.

Compatibility: every wire change is additive. Old daemons ignore the
unknown hint type and keep using the scheduled heartbeat; new daemons
against an old server simply never receive one. The cached response is
shaped exactly like a completed live discovery apart from the optional
`cached` / `cached_at` markers.

Co-authored-by: multica-agent <github@multica.ai>

* fix(agents): address review on model discovery SWR (MUL-5444)

Sol-Boy's review on #6098 found the client cache could outlive the
server's own staleness promise, and that the two changed endpoints were
still cast rather than validated.

Must-fix 1 — client freshness now derives from the served answer.
`staleTime` was a flat 5min, so a 14-minute-old snapshot (which the
server returns while queueing its own refresh) was held as fresh for
another 5min: observable staleness became server window + client window,
and the refreshed catalog never reached the tab that triggered the
refresh. `staleTime` is now a function of the query data: a `cached`
answer is stale on arrival (bound stays the server's window alone, and
the next mount/focus picks up the refreshed snapshot), while a live
discovery — which just measured the truth — is trusted for the full 5min
so a cold runtime is never re-enumerated inside one form session.
`gcTime` stays 30min, so a revisited runtime still renders from cache and
revalidates in the background; the pickers gate their spinner on
`isLoading`, which stays false throughout.

Must-fix 2 — both model-discovery responses go through a zod schema.
`POST /api/runtimes/{id}/models` and its poll companion were casting
network JSON to `RuntimeModelListRequest`, which the root CLAUDE.md
API-compatibility rules forbid. Added a lenient schema (`status` stays
`z.string()`, `supported` defaults to true, `.loose()` keeps unknown
fields) plus a fallback record whose `status` is `failed`: a malformed
body now surfaces "discovery failed" with manual entry still usable
instead of a fabricated empty catalog or an endless spinner.
`resolveRuntimeModels` was tightened to match — only an explicit
`completed` is a catalog, so an unrecognised status is an error rather
than a silent empty list, and `supported` can no longer be `undefined`.

Nit — the in-memory catalog cache now deep-copies each entry's
`Thinking` (and its level slice) and `ServiceTiers`, so it delivers the
independent value its comment promises and matches the Redis backend's
JSON round-trip semantics.

Tests: staleTime policy for cached/live/no-data; a QueryObserver test
proving the refreshed catalog reaches the same client with no blank
loading state; unknown-status and omitted-`supported` handling; schema
tests for live, cached, old-backend and nine malformed shapes; client
tests that both endpoints degrade to an explicit failure; nested-field
mutation isolation for the cache.

Co-authored-by: multica-agent <github@multica.ai>

---------

Co-authored-by: Eve <eve@multica-ai.local>
Co-authored-by: multica-agent <github@multica.ai>
2026-07-29 16:03:26 +08:00

131 lines
3.9 KiB
Go

package daemonws
import (
"log/slog"
"github.com/oklog/ulid/v2"
"github.com/multica-ai/multica/server/internal/realtime"
)
// RelayNotifier sends daemon wakeup hints to the local daemon hub and, when
// Redis is configured, publishes the same hint through the shared realtime
// relay so every API node can attempt local delivery.
type RelayNotifier struct {
local *Hub
relay realtime.RelayPublisher
}
func NewRelayNotifier(local *Hub, relay realtime.RelayPublisher) *RelayNotifier {
return &RelayNotifier{local: local, relay: relay}
}
func (n *RelayNotifier) NotifyTaskAvailable(runtimeID, taskID string) {
if runtimeID == "" {
return
}
eventID := ulid.Make().String()
if n.local != nil {
n.local.notifyTaskAvailable(runtimeID, taskID, eventID)
}
if n.relay == nil {
return
}
frame, err := taskAvailableFrame(runtimeID, taskID)
if err != nil {
M.WakeupPublishErrors.Add(1)
return
}
shardKey := taskID
if shardKey == "" {
shardKey = eventID
}
if err := n.relay.PublishWithID(realtime.ScopeDaemonRuntime, shardKey, "", frame, eventID); err != nil {
M.WakeupPublishErrors.Add(1)
slog.Warn("daemon websocket wakeup publish failed", "error", err, "runtime_id", runtimeID, "task_id", taskID)
return
}
M.WakeupPublishedTotal.Add(1)
}
func (n *RelayNotifier) NotifyRuntimeProfilesChanged(workspaceID, profileID string) {
if workspaceID == "" {
return
}
eventID := ulid.Make().String()
if n.local != nil {
n.local.notifyRuntimeProfilesChanged(workspaceID, profileID, eventID)
}
if n.relay == nil {
return
}
frame, err := runtimeProfilesChangedFrame(workspaceID, profileID)
if err != nil {
M.WakeupPublishErrors.Add(1)
return
}
if err := n.relay.PublishWithID(realtime.ScopeDaemonRuntime, workspaceID, "", frame, eventID); err != nil {
M.WakeupPublishErrors.Add(1)
slog.Warn("daemon websocket profile refresh publish failed", "error", err, "workspace_id", workspaceID, "runtime_profile_id", profileID)
return
}
M.WakeupPublishedTotal.Add(1)
}
func (n *RelayNotifier) NotifyWorkspacesChanged(userID string) {
if userID == "" {
return
}
eventID := ulid.Make().String()
if n.local != nil {
n.local.notifyWorkspacesChanged(userID, eventID)
}
if n.relay == nil {
return
}
frame, err := workspacesChangedFrame()
if err != nil {
M.WakeupPublishErrors.Add(1)
return
}
// ScopeDaemonRuntime is the relay's daemon-only transport scope; the frame
// type tells Hub.DeliverDaemonRuntime whether scopeID is a runtime,
// workspace, or user key. Keeping one transport scope preserves compatibility
// with existing relay consumers while the hub enforces user-scoped delivery.
if err := n.relay.PublishWithID(realtime.ScopeDaemonRuntime, userID, "", frame, eventID); err != nil {
M.WakeupPublishErrors.Add(1)
slog.Warn("daemon websocket workspace refresh publish failed", "error", err, "user_id", userID)
return
}
M.WakeupPublishedTotal.Add(1)
}
// NotifyPendingWork fans a runtime-scoped "heartbeat now" hint out to the local
// hub and, when Redis is configured, through the relay so the API node that
// actually holds the daemon's WebSocket delivers it (MUL-5444). Shard key is the
// runtime ID: hints for one runtime stay ordered relative to each other, and a
// dropped hint only costs the daemon its normal heartbeat delay.
func (n *RelayNotifier) NotifyPendingWork(runtimeID, kind string) {
if runtimeID == "" {
return
}
eventID := ulid.Make().String()
if n.local != nil {
n.local.notifyPendingWork(runtimeID, kind, eventID)
}
if n.relay == nil {
return
}
frame, err := pendingWorkFrame(runtimeID, kind)
if err != nil {
M.WakeupPublishErrors.Add(1)
return
}
if err := n.relay.PublishWithID(realtime.ScopeDaemonRuntime, runtimeID, "", frame, eventID); err != nil {
M.WakeupPublishErrors.Add(1)
slog.Warn("daemon websocket pending work publish failed", "error", err, "runtime_id", runtimeID, "kind", kind)
return
}
M.WakeupPublishedTotal.Add(1)
}