@@ -60,19 +60,28 @@ func (b *EventBridge) PublishEvent(agentName, userID string, evt AgentEvent) err
6060 return b .nats .Publish (subject , evt )
6161}
6262
63- // PersistObservable writes a structured observable record to the database.
63+ // PersistObservable writes a structured observable record to the database
64+ // and publishes an observable_update SSE event for real-time UI updates.
6465// The obs should be a coreTypes.Observable or compatible struct.
65- func (b * EventBridge ) PersistObservable (agentName , eventType string , obs any ) {
66+ func (b * EventBridge ) PersistObservable (agentName , userID , eventType string , obs any ) {
6667 if b .store == nil {
6768 return
6869 }
70+ payload := mustJSON (obs )
6971 b .store .AppendObservable (& AgentObservableRecord {
7072 ID : uuid .New ().String (),
71- AgentName : agentName ,
73+ AgentName : AgentKey ( userID , agentName ) ,
7274 EventType : eventType ,
73- PayloadJSON : mustJSON ( obs ) ,
75+ PayloadJSON : payload ,
7476 CreatedAt : time .Now (),
7577 })
78+ // Publish real-time SSE update (uses plain agentName for NATS subject routing)
79+ b .PublishEvent (agentName , userID , AgentEvent {
80+ AgentName : agentName ,
81+ UserID : userID ,
82+ EventType : "observable_update" ,
83+ Metadata : payload ,
84+ })
7685}
7786
7887// PublishMessage publishes a chat message event via NATS for SSE bridging.
@@ -240,7 +249,7 @@ func (b *EventBridge) handleSSEInternal(c echo.Context, agentName, userID string
240249 if evt .Metadata != "" {
241250 writeSSE (evt .EventType , evt .Metadata )
242251 }
243- case "stream_event" :
252+ case "stream_event" , "observable_update" :
244253 // Send the metadata JSON directly — React UI expects {type, content, ...}
245254 if evt .Metadata != "" {
246255 writeSSE (evt .EventType , evt .Metadata )
0 commit comments