Skip to content

Commit 75d72cf

Browse files
ptoneScion Agent (fix-interrupt-prefix)
andauthored
Support '!' prefix in broker messages as inline interrupt (#375)
* Support "!" prefix in broker messages as inline interrupt signal Messages arriving via the broker (Telegram, webhooks, direct messages) that start with "!" are now treated as interrupt messages: the "!" is stripped and the message is delivered with urgent/interrupt semantics, equivalent to --interrupt on the CLI. * Fix pointer mutation of event-bus message in interrupt prefix handling deliverToAgent mutates msg.Msg and msg.Urgent in place, but at the direct-subscriber call site the pointer comes from the event bus and is shared across all matching subscribers. Shallow-copy before mutating so the original bus message stays intact. * Handle whitespace and empty-content edge cases in interrupt prefix TrimSpace the message before checking for the "!" prefix so leading whitespace (e.g. " !restart") is handled correctly. TrimSpace the content after stripping the prefix so "! restart" works too. Default to "interrupt" when the stripped content is empty (bare "!" or "! ") to satisfy the Msg-required invariant in StructuredMessage.Validate. Adds table-driven tests covering all edge cases. * Fix CI failures: errcheck lint and no_sqlite build tag - Use t.Cleanup with explicit error discard for eventbus Close() calls in new test functions to satisfy the errcheck linter. - Add //go:build !no_sqlite constraint to signing_key_shared_test.go so newTestStore is available when `go vet -tags no_sqlite` runs. * Fix gofmt struct field alignment in edge case test --------- Co-authored-by: Scion Agent (fix-interrupt-prefix) <agent@scion.dev>
1 parent 8664cf2 commit 75d72cf

2 files changed

Lines changed: 199 additions & 0 deletions

File tree

pkg/hub/messagebroker.go

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -495,6 +495,21 @@ func (p *MessageBrokerProxy) deliverToAgent(ctx context.Context, projectID, agen
495495
return
496496
}
497497

498+
// A leading "!" in the message body acts as an inline interrupt signal:
499+
// strip the prefix and promote to urgent so the harness is interrupted
500+
// before delivery — equivalent to --interrupt on the CLI.
501+
// Shallow-copy to avoid mutating the event-bus pointer shared across subscribers.
502+
if trimmed := strings.TrimSpace(msg.Msg); strings.HasPrefix(trimmed, "!") {
503+
stripped := *msg
504+
content := strings.TrimSpace(trimmed[1:])
505+
if content == "" {
506+
content = "interrupt"
507+
}
508+
stripped.Msg = content
509+
stripped.Urgent = true
510+
msg = &stripped
511+
}
512+
498513
dispatcher := p.getDispatcher()
499514
if dispatcher == nil {
500515
p.log.Warn("No dispatcher available, cannot deliver broker message",

pkg/hub/messagebroker_test.go

Lines changed: 184 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -195,6 +195,190 @@ func TestMessageBrokerProxy_DirectMessage(t *testing.T) {
195195
}
196196
}
197197

198+
func TestMessageBrokerProxy_InterruptPrefix(t *testing.T) {
199+
s := newBrokerTestStore(t)
200+
projectID := setupBrokerTestProject(t, s)
201+
setupBrokerTestAgent(t, s, projectID, "test-agent", "running")
202+
203+
events := NewChannelEventPublisher()
204+
defer events.Close()
205+
206+
b := eventbus.NewInProcessEventBus(slog.Default())
207+
t.Cleanup(func() { _ = b.Close() })
208+
209+
dispatcher := &brokerMockDispatcher{}
210+
211+
proxy := NewMessageBrokerProxy(b, s, events, func() AgentDispatcher { return dispatcher }, slog.Default())
212+
proxy.Start()
213+
defer proxy.Stop()
214+
215+
proxy.subscribeAgent(projectID, "test-agent")
216+
217+
msg := messages.NewInstruction("user:alice", "agent:test-agent", "!restart now")
218+
if err := proxy.PublishMessage(context.Background(), projectID, msg); err != nil {
219+
t.Fatal(err)
220+
}
221+
222+
time.Sleep(100 * time.Millisecond)
223+
224+
dispatched := dispatcher.getMessages()
225+
if len(dispatched) != 1 {
226+
t.Fatalf("expected 1 dispatched message, got %d", len(dispatched))
227+
}
228+
if dispatched[0].msg != "restart now" {
229+
t.Errorf("expected message 'restart now' (! stripped), got %q", dispatched[0].msg)
230+
}
231+
if !dispatched[0].interrupt {
232+
t.Error("expected interrupt=true for !-prefixed message")
233+
}
234+
if !dispatched[0].structured.Urgent {
235+
t.Error("expected structured message Urgent=true for !-prefixed message")
236+
}
237+
}
238+
239+
func TestMessageBrokerProxy_InterruptPrefixNotStrippedWithoutBang(t *testing.T) {
240+
s := newBrokerTestStore(t)
241+
projectID := setupBrokerTestProject(t, s)
242+
setupBrokerTestAgent(t, s, projectID, "test-agent", "running")
243+
244+
events := NewChannelEventPublisher()
245+
defer events.Close()
246+
247+
b := eventbus.NewInProcessEventBus(slog.Default())
248+
t.Cleanup(func() { _ = b.Close() })
249+
250+
dispatcher := &brokerMockDispatcher{}
251+
252+
proxy := NewMessageBrokerProxy(b, s, events, func() AgentDispatcher { return dispatcher }, slog.Default())
253+
proxy.Start()
254+
defer proxy.Stop()
255+
256+
proxy.subscribeAgent(projectID, "test-agent")
257+
258+
msg := messages.NewInstruction("user:alice", "agent:test-agent", "hello agent")
259+
if err := proxy.PublishMessage(context.Background(), projectID, msg); err != nil {
260+
t.Fatal(err)
261+
}
262+
263+
time.Sleep(100 * time.Millisecond)
264+
265+
dispatched := dispatcher.getMessages()
266+
if len(dispatched) != 1 {
267+
t.Fatalf("expected 1 dispatched message, got %d", len(dispatched))
268+
}
269+
if dispatched[0].msg != "hello agent" {
270+
t.Errorf("expected message 'hello agent' unchanged, got %q", dispatched[0].msg)
271+
}
272+
if dispatched[0].interrupt {
273+
t.Error("expected interrupt=false for non-!-prefixed message")
274+
}
275+
}
276+
277+
func TestMessageBrokerProxy_InterruptPrefixEdgeCases(t *testing.T) {
278+
tests := []struct {
279+
name string
280+
input string
281+
wantMsg string
282+
wantInterrupt bool
283+
wantUrgent bool
284+
}{
285+
{"bare bang", "!", "interrupt", true, true},
286+
{"bang with trailing spaces", "! ", "interrupt", true, true},
287+
{"leading whitespace before bang", " !restart", "restart", true, true},
288+
{"whitespace between bang and content", "! restart", "restart", true, true},
289+
{"leading and inner whitespace", " ! restart now ", "restart now", true, true},
290+
{"normal message no prefix", "hello", "hello", false, false},
291+
}
292+
293+
for _, tt := range tests {
294+
t.Run(tt.name, func(t *testing.T) {
295+
s := newBrokerTestStore(t)
296+
projectID := setupBrokerTestProject(t, s)
297+
setupBrokerTestAgent(t, s, projectID, "test-agent", "running")
298+
299+
events := NewChannelEventPublisher()
300+
defer events.Close()
301+
302+
bus := eventbus.NewInProcessEventBus(slog.Default())
303+
t.Cleanup(func() { _ = bus.Close() })
304+
305+
dispatcher := &brokerMockDispatcher{}
306+
307+
proxy := NewMessageBrokerProxy(bus, s, events, func() AgentDispatcher { return dispatcher }, slog.Default())
308+
proxy.Start()
309+
defer proxy.Stop()
310+
311+
proxy.subscribeAgent(projectID, "test-agent")
312+
313+
msg := messages.NewInstruction("user:alice", "agent:test-agent", tt.input)
314+
if err := proxy.PublishMessage(context.Background(), projectID, msg); err != nil {
315+
t.Fatal(err)
316+
}
317+
318+
time.Sleep(100 * time.Millisecond)
319+
320+
dispatched := dispatcher.getMessages()
321+
if len(dispatched) != 1 {
322+
t.Fatalf("expected 1 dispatched message, got %d", len(dispatched))
323+
}
324+
if dispatched[0].msg != tt.wantMsg {
325+
t.Errorf("msg = %q, want %q", dispatched[0].msg, tt.wantMsg)
326+
}
327+
if dispatched[0].interrupt != tt.wantInterrupt {
328+
t.Errorf("interrupt = %v, want %v", dispatched[0].interrupt, tt.wantInterrupt)
329+
}
330+
if dispatched[0].structured.Urgent != tt.wantUrgent {
331+
t.Errorf("Urgent = %v, want %v", dispatched[0].structured.Urgent, tt.wantUrgent)
332+
}
333+
})
334+
}
335+
}
336+
337+
func TestMessageBrokerProxy_InterruptPrefixPersistence(t *testing.T) {
338+
s := newBrokerTestStore(t)
339+
projectID := setupBrokerTestProject(t, s)
340+
agent := setupBrokerTestAgent(t, s, projectID, "persist-agent", "running")
341+
342+
events := NewChannelEventPublisher()
343+
defer events.Close()
344+
345+
b := eventbus.NewInProcessEventBus(slog.Default())
346+
t.Cleanup(func() { _ = b.Close() })
347+
348+
dispatcher := &brokerMockDispatcher{}
349+
350+
proxy := NewMessageBrokerProxy(b, s, events, func() AgentDispatcher { return dispatcher }, slog.Default())
351+
proxy.Start()
352+
defer proxy.Stop()
353+
354+
proxy.subscribeAgent(projectID, "persist-agent")
355+
356+
msg := messages.NewInstruction("user:alice", "agent:persist-agent", "!urgent task")
357+
msg.SenderID = "user-alice-id"
358+
msg.RecipientID = agent.ID
359+
if err := proxy.PublishMessage(context.Background(), projectID, msg); err != nil {
360+
t.Fatal(err)
361+
}
362+
363+
time.Sleep(100 * time.Millisecond)
364+
365+
// Verify the persisted message has the stripped content and urgent flag
366+
ctx := context.Background()
367+
result, err := s.ListMessages(ctx, store.MessageFilter{AgentID: agent.ID}, store.ListOptions{})
368+
if err != nil {
369+
t.Fatalf("failed to list messages: %v", err)
370+
}
371+
if len(result.Items) != 1 {
372+
t.Fatalf("expected 1 persisted message, got %d", len(result.Items))
373+
}
374+
if result.Items[0].Msg != "urgent task" {
375+
t.Errorf("expected persisted msg 'urgent task', got %q", result.Items[0].Msg)
376+
}
377+
if !result.Items[0].Urgent {
378+
t.Error("expected persisted message Urgent=true")
379+
}
380+
}
381+
198382
func TestMessageBrokerProxy_ProjectBroadcast(t *testing.T) {
199383
s := newBrokerTestStore(t)
200384
projectID := setupBrokerTestProject(t, s)

0 commit comments

Comments
 (0)