diff --git a/internal/bridge/daimon/AGENTS.md b/internal/bridge/daimon/AGENTS.md index bf889c1..e3634dd 100644 --- a/internal/bridge/daimon/AGENTS.md +++ b/internal/bridge/daimon/AGENTS.md @@ -10,6 +10,11 @@ authenticated v2 durable wake-acceptance endpoint. - ACK only after a matching durable acceptance and local receipt job are persisted. Follow cognition asynchronously, publishing terminal reply text with an idempotent Moltnet message id. +- A strict `noopolis.daimon.work-blocked.v1` descriptor nested in a 409 + stopped response is backpressure, not wake failure. Preserve the cursor and + delivery identity; retry after its delay clamped to 1 second–5 minutes. + Ordinary budget pauses still accept durable ownership through unchanged 202 + receipts. Never interpret arbitrary 409 responses as accepted or deferred. - An explicit `moltnet send` made during a Daimon wake and the terminal receipt fallback share one target-scoped idempotent publication slot. The first durable message wins; terminal-only agents still publish through the fallback. diff --git a/internal/bridge/daimon/blocked.go b/internal/bridge/daimon/blocked.go new file mode 100644 index 0000000..fafb6ce --- /dev/null +++ b/internal/bridge/daimon/blocked.go @@ -0,0 +1,65 @@ +package daimon + +import ( + "bytes" + "encoding/json" + "fmt" + "net/http" + "time" + + "github.com/noopolis/moltnet/internal/bridge/loop" +) + +const workBlockedVersion = "noopolis.daimon.work-blocked.v1" + +// Only the versioned backpressure response is a deferral. Legacy 409 errors +// (unknown agent, conflicting delivery, malformed request) remain failures. +func decodeBlockedResponse(response *http.Response) error { + failure := fmt.Errorf("control url returned %s", response.Status) + fields, err := decodeExactObject(response.Body) + if err != nil || requireExactFields(fields, "version", "state", "code", "blocked") != nil { + return failure + } + version, err := requiredString(fields, "version") + if err != nil || version != wakeAcceptanceVersion { + return failure + } + state, err := requiredString(fields, "state") + if err != nil || state != "stopped" { + return failure + } + code, err := requiredString(fields, "code") + if err != nil || (code != "host_stopping" && code != "host_stopped") { + return failure + } + blocked, err := decodeExactObject(bytes.NewReader(fields["blocked"])) + if err != nil || requireExactFields(blocked, "version", "reason", "retry_after_ms") != nil { + return failure + } + version, err = requiredString(blocked, "version") + if err != nil || version != workBlockedVersion { + return failure + } + reason, err := requiredString(blocked, "reason") + if err != nil || !validBlockedReason(reason) { + return failure + } + var retryMS int64 + if err := json.Unmarshal(blocked["retry_after_ms"], &retryMS); err != nil || retryMS <= 0 { + return failure + } + // Clamp before converting to Duration, which could otherwise overflow. + if retryMS > int64((5*time.Minute)/time.Millisecond) { + retryMS = int64((5 * time.Minute) / time.Millisecond) + } + return &loop.ControlDeferredError{Reason: reason, RetryAfter: time.Duration(retryMS) * time.Millisecond} +} + +func validBlockedReason(reason string) bool { + switch reason { + case "operator_stop", "ledger_unavailable", "host_stopping", "host_stopped", "queue_full": + return true + default: + return false + } +} diff --git a/internal/bridge/daimon/blocked_adapter_test.go b/internal/bridge/daimon/blocked_adapter_test.go new file mode 100644 index 0000000..e6750d2 --- /dev/null +++ b/internal/bridge/daimon/blocked_adapter_test.go @@ -0,0 +1,175 @@ +package daimon + +import ( + "context" + "encoding/json" + "fmt" + "io" + "net/http" + "net/http/httptest" + "path/filepath" + "sync" + "testing" + "time" + + "github.com/gorilla/websocket" + "github.com/noopolis/moltnet/pkg/bridgeconfig" + "github.com/noopolis/moltnet/pkg/protocol" +) + +func TestAdapterDefersWithoutACKThenRecoversSameDelivery(t *testing.T) { + t.Setenv("DAIMON_ADAPTER_TOKEN", "test-bearer") + ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) + defer cancel() + delivery := daimonDelivery() + delivery.Message = "[room research] writer\nhello" + event := protocol.Event{ + ID: "cursor_pending", Type: protocol.EventTypeMessageCreated, NetworkID: "local", + Message: &protocol.Message{ + ID: "msg_1", NetworkID: "local", Target: delivery.Target, + From: protocol.Actor{Type: "agent", ID: "writer"}, Mentions: []string{"researcher"}, + Parts: []protocol.Part{{Kind: protocol.PartKindText, Text: "hello"}}, CreatedAt: delivery.OccurredAt, + }, CreatedAt: delivery.OccurredAt, + } + config := daimonConfig("http://control.invalid") + var mu sync.Mutex + var requestBodies []string + var requestTimes []time.Time + var attachments, failures, publications int + acked := make(chan struct{}) + published := make(chan struct{}) + controlServer := httptest.NewServer(http.HandlerFunc(func(response http.ResponseWriter, request *http.Request) { + switch request.URL.Path { + case "/v2/wakes": + body, _ := io.ReadAll(request.Body) + mu.Lock() + requestBodies = append(requestBodies, string(body)) + requestTimes = append(requestTimes, time.Now()) + attempt := len(requestBodies) + mu.Unlock() + if attempt <= 2 { + response.WriteHeader(http.StatusConflict) + _, _ = response.Write([]byte(blockedBody("operator_stop", 1000))) + return + } + response.WriteHeader(http.StatusAccepted) + _, _ = response.Write([]byte(acceptanceBody(t, config, delivery))) + case "/v2/wake-receipts/11111111-1111-4111-8111-111111111111": + select { + case <-acked: + case <-request.Context().Done(): + return + } + _, _ = response.Write([]byte(receiptBody(t, config, delivery, "completed", "recovered reply"))) + default: + response.WriteHeader(http.StatusNotFound) + } + })) + defer controlServer.Close() + upgrader := websocket.Upgrader{CheckOrigin: func(*http.Request) bool { return true }} + moltnetServer := httptest.NewServer(http.HandlerFunc(func(response http.ResponseWriter, request *http.Request) { + switch request.URL.Path { + case "/v1/attach": + mu.Lock() + attachments++ + attempt := attachments + mu.Unlock() + conn, err := upgrader.Upgrade(response, request, nil) + if err != nil { + t.Error(err) + return + } + defer conn.Close() + _ = conn.SetReadDeadline(time.Now().Add(5 * time.Second)) + identify := blockedHandshake(t, conn) + if attempt == 1 { + // A skipped, unrelated event establishes the prior ACK cursor. + prior := protocol.Event{ID: "cursor_saved", Type: "unrelated", NetworkID: "local"} + _ = conn.WriteJSON(protocol.AttachmentFrame{Op: protocol.AttachmentOpEvent, Version: protocol.AttachmentProtocolV1, Cursor: prior.ID, Event: &prior}) + var ack protocol.AttachmentFrame + if err := conn.ReadJSON(&ack); err != nil || ack.Op != protocol.AttachmentOpAck || ack.Cursor != prior.ID { + t.Errorf("prior cursor not ACKed: %#v, %v", ack, err) + return + } + } else if identify.Cursor != "cursor_saved" { + t.Errorf("deferred delivery advanced cursor: %q", identify.Cursor) + } + _ = conn.WriteJSON(protocol.AttachmentFrame{Op: protocol.AttachmentOpEvent, Version: protocol.AttachmentProtocolV1, Cursor: event.ID, Event: &event}) + var frame protocol.AttachmentFrame + err = conn.ReadJSON(&frame) + if attempt <= 2 { + if err == nil { + t.Errorf("blocked delivery sent ACK/error frame: %#v", frame) + } + return + } + if err != nil || frame.Op != protocol.AttachmentOpAck || frame.Cursor != event.ID { + t.Errorf("accepted delivery failed ACK: %#v, %v", frame, err) + return + } + close(acked) + select { + case <-published: + case <-ctx.Done(): + t.Error("recovered receipt did not publish") + return + } + _ = conn.WriteControl(websocket.CloseMessage, websocket.FormatCloseMessage(websocket.CloseNormalClosure, "done"), time.Now().Add(time.Second)) + case "/v1/messages": + var payload protocol.SendMessageRequest + if err := json.NewDecoder(request.Body).Decode(&payload); err != nil { + t.Error(err) + return + } + mu.Lock() + publications++ + mu.Unlock() + if len(payload.Parts) != 1 || payload.Parts[0].Text != "recovered reply" { + t.Errorf("unexpected recovered reply: %#v", payload.Parts) + } + _, _ = fmt.Fprintf(response, `{"message_id":%q,"event_id":"evt_reply","accepted":true}`, payload.ID) + close(published) + case "/v1/agents/wake-failed": + mu.Lock() + failures++ + mu.Unlock() + response.WriteHeader(http.StatusAccepted) + default: + response.WriteHeader(http.StatusNotFound) + } + })) + defer moltnetServer.Close() + runConfig := config + runConfig.Moltnet.BaseURL = moltnetServer.URL + runConfig.Runtime.ControlURL = controlServer.URL + runConfig.Runtime.TokenEnv = "DAIMON_ADAPTER_TOKEN" + runConfig.Runtime.ReceiptStorePath = filepath.Join(t.TempDir(), "private", "receipts.json") + runConfig.Rooms = []bridgeconfig.RoomBinding{{ID: "research", Wake: bridgeconfig.WakeMentions}} + if err := New().Run(ctx, runConfig); err != nil { + t.Fatal(err) + } + mu.Lock() + defer mu.Unlock() + if len(requestBodies) != 3 || attachments != 3 || failures != 0 || publications != 1 { + t.Fatalf("requests=%d attachments=%d failures=%d publications=%d, want 3/3/0/1", len(requestBodies), attachments, failures, publications) + } + for i := 1; i < len(requestBodies); i++ { + if requestBodies[i] != requestBodies[0] { + t.Fatal("retry changed delivery identity or wake content") + } + if requestTimes[i].Sub(requestTimes[i-1]) < time.Second { + t.Fatal("retried before the runtime's deferred interval") + } + } +} + +func blockedHandshake(t *testing.T, conn *websocket.Conn) protocol.AttachmentFrame { + t.Helper() + _ = conn.WriteJSON(protocol.AttachmentFrame{Op: protocol.AttachmentOpHello, Version: protocol.AttachmentProtocolV1, HeartbeatIntervalMS: 30000}) + var identify protocol.AttachmentFrame + if err := conn.ReadJSON(&identify); err != nil { + t.Error(err) + } + _ = conn.WriteJSON(protocol.AttachmentFrame{Op: protocol.AttachmentOpReady, Version: protocol.AttachmentProtocolV1, NetworkID: "local", AgentID: "researcher"}) + return identify +} diff --git a/internal/bridge/daimon/blocked_test.go b/internal/bridge/daimon/blocked_test.go new file mode 100644 index 0000000..30f4daf --- /dev/null +++ b/internal/bridge/daimon/blocked_test.go @@ -0,0 +1,71 @@ +package daimon + +import ( + "errors" + "fmt" + "net/http" + "strings" + "testing" + "time" + + "github.com/noopolis/moltnet/internal/bridge/loop" +) + +func blockedBody(reason string, retryMS int64) string { + return fmt.Sprintf(`{"version":"noopolis.daimon.wake-acceptance.v2","state":"stopped","code":"host_stopping","blocked":{"version":"noopolis.daimon.work-blocked.v1","reason":%q,"retry_after_ms":%d}}`, reason, retryMS) +} + +func TestCodecRecognizesOnlyVersionedDeferral(t *testing.T) { + for _, reason := range []string{"operator_stop", "ledger_unavailable", "host_stopping", "host_stopped", "queue_full"} { + t.Run(reason, func(t *testing.T) { + result, err := NewCodec("").DecodeResponse(daimonConfig("http://control.invalid"), daimonDelivery(), daimonResponse(http.StatusConflict, blockedBody(reason, 30000))) + var deferred *loop.ControlDeferredError + if !errors.As(err, &deferred) || deferred.Reason != reason || deferred.RetryAfter != 30*time.Second { + t.Fatalf("expected typed deferral, got %#v, %v", deferred, err) + } + if result.Acceptance != nil || result.Publish { + t.Fatalf("deferred work must not ACK or publish: %#v", result) + } + }) + } +} + +func TestCodecRejectsMalformedBlockedResponsesWithoutLeakingMaterial(t *testing.T) { + valid := blockedBody("operator_stop", 30000) + for name, body := range map[string]string{ + "legacy stop": `{"version":"noopolis.daimon.wake-acceptance.v2","state":"stopped","code":"host_stopping"}`, + "legacy rejection": `{"version":"noopolis.daimon.wake-acceptance.v2","state":"rejected","code":"unknown_agent"}`, + "wrong state": strings.Replace(valid, `"state":"stopped"`, `"state":"accepted"`, 1), + "wrong code": strings.Replace(valid, `"code":"host_stopping"`, `"code":"unknown_agent"`, 1), + "wrong outer version": strings.Replace(valid, wakeAcceptanceVersion, "unknown", 1), + "wrong nested version": strings.Replace(valid, workBlockedVersion, "unknown", 1), + "unknown reason": blockedBody("response-canary", 30000), + "negative delay": blockedBody("operator_stop", -1), + "zero delay": blockedBody("operator_stop", 0), + "fractional delay": strings.Replace(valid, "30000", "0.1", 1), + "null delay": strings.Replace(valid, "30000", "null", 1), + "string delay": strings.Replace(valid, "30000", `"30000"`, 1), + "missing delay": strings.Replace(valid, `,"retry_after_ms":30000`, "", 1), + "extra nested field": strings.Replace(valid, `"reason":`, `"extra":true,"reason":`, 1), + "duplicate nested field": strings.Replace(valid, `"reason":`, `"reason":"operator_stop","reason":`, 1), + "extra outer field": strings.Replace(valid, `"state":`, `"extra":true,"state":`, 1), + "duplicate outer field": strings.Replace(valid, `"state":`, `"state":"stopped","state":`, 1), + "trailing body": valid + `{}`, + } { + t.Run(name, func(t *testing.T) { + result, err := NewCodec("").DecodeResponse(daimonConfig("http://control.invalid"), daimonDelivery(), daimonResponse(http.StatusConflict, body)) + var deferred *loop.ControlDeferredError + if err == nil || errors.As(err, &deferred) || result.Acceptance != nil || result.Publish || strings.Contains(err.Error(), "response-canary") { + t.Fatalf("malformed response accepted or leaked material: %#v, %v", result, err) + } + }) + } +} + +func TestCodecClampsDeferredDelayBeforeDurationConversion(t *testing.T) { + _, err := NewCodec("").DecodeResponse(daimonConfig("http://control.invalid"), daimonDelivery(), daimonResponse(http.StatusConflict, blockedBody("operator_stop", 9223372036854775807))) + var deferred *loop.ControlDeferredError + if !errors.As(err, &deferred) || deferred.RetryAfter != 5*time.Minute { + t.Fatalf("expected bounded delay, got %#v, %v", deferred, err) + } +} diff --git a/internal/bridge/daimon/codec.go b/internal/bridge/daimon/codec.go index ca2ce7b..cefce2d 100644 --- a/internal/bridge/daimon/codec.go +++ b/internal/bridge/daimon/codec.go @@ -98,6 +98,9 @@ func (*Codec) DecodeResponse( delivery loop.ControlDelivery, response *http.Response, ) (loop.ControlResult, error) { + if response.StatusCode == http.StatusConflict { + return loop.ControlResult{}, decodeBlockedResponse(response) + } if response.StatusCode != http.StatusAccepted { return loop.ControlResult{}, fmt.Errorf("control url returned %s", response.Status) } diff --git a/internal/bridge/loop/control.go b/internal/bridge/loop/control.go index 017bf8a..3a019d9 100644 --- a/internal/bridge/loop/control.go +++ b/internal/bridge/loop/control.go @@ -102,7 +102,13 @@ func RunControlLoopWithCodec(ctx context.Context, config bridgeconfig.Config, co } } - if err == nil || ctx.Err() != nil { + // Cancellation can race with a deferred delivery unwinding the + // stream. Shutdown has the same clean result as the wait branch + // below; the last retryable delivery error is no longer actionable. + if ctx.Err() != nil { + return nil + } + if err == nil { return err } attempt++ @@ -110,7 +116,7 @@ func RunControlLoopWithCodec(ctx context.Context, config bridgeconfig.Config, co select { case <-ctx.Done(): return nil - case <-time.After(backoff.Delay(attempt)): + case <-time.After(controlReconnectDelay(err, backoff.Delay(attempt))): } } } diff --git a/internal/bridge/loop/control_deferred.go b/internal/bridge/loop/control_deferred.go new file mode 100644 index 0000000..bfea082 --- /dev/null +++ b/internal/bridge/loop/control_deferred.go @@ -0,0 +1,43 @@ +package loop + +import ( + "errors" + "time" +) + +const ( + minControlDeferredDelay = time.Second + maxControlDeferredDelay = 5 * time.Minute +) + +// ControlDeferredError means the runtime has not accepted ownership yet, but +// may accept the same delivery later. It is backpressure, not a failed wake: +// leave the attachment cursor unchanged, then retry after a bounded delay. +// Runtime codecs validate the wire reason before constructing this error. +type ControlDeferredError struct { + Reason string + RetryAfter time.Duration +} + +func (e *ControlDeferredError) Error() string { return "runtime work deferred: " + e.Reason } + +func controlDeferral(err error) (*ControlDeferredError, bool) { + var deferred *ControlDeferredError + ok := errors.As(err, &deferred) + return deferred, ok +} + +func controlReconnectDelay(err error, fallback time.Duration) time.Duration { + deferred, ok := controlDeferral(err) + if !ok { + return fallback + } + delay := deferred.RetryAfter + if delay < minControlDeferredDelay { + return minControlDeferredDelay + } + if delay > maxControlDeferredDelay { + return maxControlDeferredDelay + } + return delay +} diff --git a/internal/bridge/loop/control_deferred_test.go b/internal/bridge/loop/control_deferred_test.go new file mode 100644 index 0000000..3c77b97 --- /dev/null +++ b/internal/bridge/loop/control_deferred_test.go @@ -0,0 +1,119 @@ +package loop + +import ( + "context" + "fmt" + "net/http" + "testing" + "time" + + "github.com/gorilla/websocket" + "github.com/noopolis/moltnet/pkg/bridgeconfig" + "github.com/noopolis/moltnet/pkg/protocol" +) + +type deferredTestCodec struct { + *legacyControlCodec + onDeferred func() +} + +func (c *deferredTestCodec) DecodeResponse(bridgeconfig.Config, ControlDelivery, *http.Response) (ControlResult, error) { + if c.onDeferred != nil { + c.onDeferred() + } + return ControlResult{}, &ControlDeferredError{Reason: "operator_stop", RetryAfter: time.Minute} +} + +func TestDeferredDeliveryDoesNotConsumeFailureBudgetOrReportFailure(t *testing.T) { + harness := newControlRetryTestHarness(t, http.StatusConflict, func(*testing.T, *websocket.Conn, int) {}) + config := harness.config() + deliveries := newControlDeliveryTracker() + for attempt := 0; attempt < maxControlDeliveryAttempts+2; attempt++ { + err := deliverControlMessage(t.Context(), http.DefaultClient, NewMoltnetClient(config), config, + &deferredTestCodec{legacyControlCodec: &legacyControlCodec{}}, *permanentlyFailingEvent("evt_deferred"), deliveries) + if _, ok := controlDeferral(err); !ok { + t.Fatalf("expected non-ACKing deferral, got %v", err) + } + } + _, requests, reports := harness.counts() + state := deliveries.stateFor("evt_deferred") + if requests != maxControlDeliveryAttempts+2 || reports != 0 || state.attempts != 0 || state.permanent { + t.Fatalf("deferral lost retryability or reported failure: requests=%d reports=%d state=%#v", requests, reports, state) + } +} + +func TestControlDeferralReconnectDelayIsBounded(t *testing.T) { + for _, test := range []struct{ input, want time.Duration }{ + {0, time.Second}, {-time.Second, time.Second}, {time.Millisecond, time.Second}, + {30 * time.Second, 30 * time.Second}, {time.Hour, 5 * time.Minute}, + } { + err := fmt.Errorf("wrapped: %w", &ControlDeferredError{Reason: "operator_stop", RetryAfter: test.input}) + if got := controlReconnectDelay(err, time.Millisecond); got != test.want { + t.Fatalf("delay(%s)=%s, want %s", test.input, got, test.want) + } + } + if got := controlReconnectDelay(fmt.Errorf("legacy error"), time.Millisecond); got != time.Millisecond { + t.Fatalf("changed legacy delay: %s", got) + } +} + +func TestControlDeferralDoesNotSendAttachmentError(t *testing.T) { + err := &ControlDeferredError{Reason: "operator_stop", RetryAfter: time.Minute} + got := reportAttachmentHandlerError(func(protocol.AttachmentFrame) error { + t.Fatal("backpressure published an attachment error") + return nil + }, err) + if got != err { + t.Fatalf("lost replay signal: %v", got) + } +} + +func TestDeferredReconnectWaitCanBeCancelled(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + harness := newControlRetryTestHarness(t, http.StatusConflict, func(t *testing.T, conn *websocket.Conn, attempt int) { + if attempt > 1 { + t.Error("retried before the deferred delay expired") + } + attachHandshake(t, conn) + event := permanentlyFailingEvent("evt_deferred") + _ = conn.WriteJSON(protocol.AttachmentFrame{Op: protocol.AttachmentOpEvent, Version: protocol.AttachmentProtocolV1, Cursor: event.ID, Event: event}) + var frame protocol.AttachmentFrame + if err := conn.ReadJSON(&frame); err == nil { + t.Errorf("deferred delivery emitted frame %#v", frame) + } + cancel() + }) + started := time.Now() + if err := RunControlLoopWithCodec(ctx, harness.config(), &deferredTestCodec{legacyControlCodec: &legacyControlCodec{}}); err != nil { + t.Fatal(err) + } + if time.Since(started) > time.Second { + t.Fatal("shutdown waited out the deferred retry delay") + } + _, requests, reports := harness.counts() + if requests != 1 || reports != 0 { + t.Fatalf("requests=%d reports=%d, want 1/0", requests, reports) + } +} + +func TestDeferralCancelledBeforeStreamUnwindsReturnsCleanly(t *testing.T) { + ctx, cancel := context.WithCancel(t.Context()) + defer cancel() + harness := newControlRetryTestHarness(t, http.StatusConflict, func(t *testing.T, conn *websocket.Conn, _ int) { + attachHandshake(t, conn) + event := permanentlyFailingEvent("evt_deferred") + _ = conn.WriteJSON(protocol.AttachmentFrame{Op: protocol.AttachmentOpEvent, Version: protocol.AttachmentProtocolV1, Cursor: event.ID, Event: event}) + var frame protocol.AttachmentFrame + if err := conn.ReadJSON(&frame); err == nil { + t.Errorf("cancelled deferral emitted frame %#v", frame) + } + }) + codec := &deferredTestCodec{legacyControlCodec: &legacyControlCodec{}, onDeferred: cancel} + if err := RunControlLoopWithCodec(ctx, harness.config(), codec); err != nil { + t.Fatalf("cancelled loop returned prior delivery error: %v", err) + } + _, requests, reports := harness.counts() + if requests != 1 || reports != 0 { + t.Fatalf("requests=%d reports=%d, want 1/0", requests, reports) + } +} diff --git a/internal/bridge/loop/control_delivery.go b/internal/bridge/loop/control_delivery.go index 68a5372..590d37d 100644 --- a/internal/bridge/loop/control_delivery.go +++ b/internal/bridge/loop/control_delivery.go @@ -82,6 +82,12 @@ func deliverControlMessage( } return publishErr } + if _, deferred := controlDeferral(err); deferred { + // Backpressure neither consumes the delivery failure budget nor + // authorizes an ACK. Reconnect at the runtime's bounded cadence. + state.attempts-- + return err + } switch classifyControlError(ctx, err) { case controlErrorFatal: diff --git a/internal/bridge/loop/moltnet.go b/internal/bridge/loop/moltnet.go index f00145b..247b00e 100644 --- a/internal/bridge/loop/moltnet.go +++ b/internal/bridge/loop/moltnet.go @@ -169,6 +169,11 @@ func reportAttachmentHandlerError(write func(protocol.AttachmentFrame) error, er if err == nil { return nil } + if _, deferred := controlDeferral(err); deferred { + // A deferred delivery remains unACKed for replay. It is an expected + // runtime availability state, not an attachment handler failure. + return err + } message := strings.TrimSpace(err.Error()) if message == "" { message = "bridge handler failed"