diff --git a/event/stream/decode.go b/event/stream/decode.go index 2559a8b..ac3b88e 100644 --- a/event/stream/decode.go +++ b/event/stream/decode.go @@ -7,8 +7,11 @@ import ( "github.com/kerberos-io/onvif/event" ) -// Decode converts a single ONVIF NotificationMessage into the package's -// normalized Event representation. +// decode converts a single ONVIF NotificationMessage into the package's +// normalized Event representation. Unexported because the only intended +// caller is the Stream; downstream consumers receive decoded Events on +// the Events channel. Tests reach decode directly because they're in +// the same package. // // deviceID is supplied by the caller because the message itself does not // identify the originating camera. observedAt is recorded verbatim as @@ -18,7 +21,7 @@ import ( // When the Topic does not match any classifier rule the returned Event // has Kind == KindUnknown but Source, Data and Topic are still populated // so consumers can fall back to inspecting the wire form. -func Decode(msg event.NotificationMessage, deviceID string, observedAt time.Time) Event { +func decode(msg event.NotificationMessage, deviceID string, observedAt time.Time) Event { topic := string(msg.Topic.TopicKinds) desc := msg.Message.Message return Event{ diff --git a/event/stream/decode_test.go b/event/stream/decode_test.go index edce4c8..2d52800 100644 --- a/event/stream/decode_test.go +++ b/event/stream/decode_test.go @@ -52,7 +52,7 @@ func TestDecode_MotionActive(t *testing.T) { map[string]string{"IsMotion": "true"}, ) - ev := Decode(in, "axis-cam-01", observedAt) + ev := decode(in, "axis-cam-01", observedAt) assert.Equal(t, KindMotion, ev.Kind) assert.Equal(t, StateActive, ev.State) @@ -74,7 +74,7 @@ func TestDecode_MotionInactive(t *testing.T) { nil, map[string]string{"State": "false"}, ) - ev := Decode(in, "dev", time.Now()) + ev := decode(in, "dev", time.Now()) assert.Equal(t, KindMotion, ev.Kind) assert.Equal(t, StateInactive, ev.State) } @@ -88,7 +88,7 @@ func TestDecode_HanwhaNumericMotionValue(t *testing.T) { nil, map[string]string{"Motion": "1"}, ) - ev := Decode(in, "dev", time.Now()) + ev := decode(in, "dev", time.Now()) assert.Equal(t, KindMotion, ev.Kind) assert.Equal(t, StateActive, ev.State) } @@ -103,7 +103,7 @@ func TestDecode_AvigilonActiveLiteral(t *testing.T) { map[string]string{"RelayToken": "Relay-1"}, map[string]string{"LogicalState": "active"}, ) - ev := Decode(in, "dev", time.Now()) + ev := decode(in, "dev", time.Now()) assert.Equal(t, KindDigitalOutput, ev.Kind) assert.Equal(t, StateActive, ev.State) assert.Equal(t, "Relay-1", ev.Source["RelayToken"]) @@ -124,7 +124,7 @@ func TestDecode_AxisObjectAnalyticsMultiItem(t *testing.T) { "confidence": "92", }, ) - ev := Decode(in, "dev", time.Now()) + ev := decode(in, "dev", time.Now()) assert.Equal(t, KindObjectDetected, ev.Kind) assert.Equal(t, StateActive, ev.State) assert.Equal(t, "Human", ev.Data["classType"]) @@ -142,7 +142,7 @@ func TestDecode_LineDetectorCrossedHasNoState(t *testing.T) { map[string]string{"VideoSourceConfigurationToken": "vsct0", "Rule": "LineRule"}, map[string]string{"ObjectId": "42"}, ) - ev := Decode(in, "dev", time.Now()) + ev := decode(in, "dev", time.Now()) assert.Equal(t, KindObjectDetected, ev.Kind) assert.Equal(t, StateUnknown, ev.State) assert.Equal(t, "42", ev.Data["ObjectId"]) @@ -158,7 +158,7 @@ func TestDecode_UnknownTopicStillPreservesWireData(t *testing.T) { nil, map[string]string{"Custom": "true"}, ) - ev := Decode(in, "dev", time.Now()) + ev := decode(in, "dev", time.Now()) assert.Equal(t, KindUnknown, ev.Kind) assert.Equal(t, "tns1:UserAlarm/IVA", ev.Topic) assert.Equal(t, "true", ev.Data["Custom"]) @@ -179,7 +179,7 @@ func TestDecode_PropertyOperationVariants(t *testing.T) { for _, tc := range tests { t.Run(tc.name, func(t *testing.T) { in := msg("tns1:VideoSource/MotionAlarm", tc.in, "", nil, nil) - ev := Decode(in, "dev", time.Now()) + ev := decode(in, "dev", time.Now()) assert.Equal(t, tc.want, ev.Operation) }) } @@ -200,7 +200,7 @@ func TestDecode_DeviceTimeParsing(t *testing.T) { for _, tc := range tests { t.Run(tc.name, func(t *testing.T) { in := msg("tns1:VideoSource/MotionAlarm", "Changed", tc.in, nil, nil) - ev := Decode(in, "dev", time.Now()) + ev := decode(in, "dev", time.Now()) if tc.want.IsZero() { assert.True(t, ev.DeviceTime.IsZero(), "DeviceTime=%v", ev.DeviceTime) } else { @@ -215,7 +215,7 @@ func TestDecode_EmptySourceAndDataYieldNilMaps(t *testing.T) { // safely len() and index into Source/Data without nil-checking, but // we do not allocate an empty map for empty notifications. in := msg("tns1:VideoSource/MotionAlarm", "Changed", "", nil, nil) - ev := Decode(in, "dev", time.Now()) + ev := decode(in, "dev", time.Now()) assert.Nil(t, ev.Source) assert.Nil(t, ev.Data) } @@ -242,7 +242,7 @@ func TestDecode_StateValueIsCaseInsensitive(t *testing.T) { t.Run(tc.name, func(t *testing.T) { in := msg("tns1:VideoSource/MotionAlarm", "Changed", "", nil, map[string]string{"State": tc.value}) - ev := Decode(in, "dev", time.Now()) + ev := decode(in, "dev", time.Now()) assert.Equal(t, tc.want, ev.State) }) } diff --git a/event/stream/doc.go b/event/stream/doc.go index 0223ae1..8121525 100644 --- a/event/stream/doc.go +++ b/event/stream/doc.go @@ -1,13 +1,59 @@ -// Package stream will provide a long-running, channel-based consumer for -// ONVIF device events. It is meant to hide the SOAP/XML, pull-point -// subscription lifecycle, subscription renewal and vendor-specific topic -// conventions behind a typed Event stream. +// Package stream is a typed, channel-based consumer for ONVIF device +// events. It hides the SOAP/XML, pull-point subscription lifecycle, +// subscription renewal and vendor-specific topic conventions behind a +// single Event stream. // -// This file lays down the value types (Kind, State, PropertyOperation, -// Event) and the topic Classifier. The Stream type, its NewStream -// constructor and the Events/Errors channels land in follow-up changes. +// # Usage // -// The package classifies vendor-specific topic strings (AXIS, Hikvision, -// Avigilon, Hanwha, Bosch, Dahua) into a small set of normalized Kind -// values so callers do not need to special-case device manufacturers. +// dev, _ := onvif.NewDevice(onvif.DeviceParams{Xaddr: "...", Username: "...", Password: "..."}) +// s, err := stream.NewStream(ctx, dev, stream.Options{DeviceID: "front-door"}) +// if err != nil { /* construction failed: auth, network, or camera does not advertise events */ } +// defer s.Close() +// +// for ev := range s.Events() { +// switch ev.Kind { +// case stream.KindMotion: +// if ev.State == stream.StateActive { /* start recording */ } +// } +// } +// +// # Invariants +// +// NewStream performs network I/O. It returns once the +// CreatePullPointSubscription call has succeeded; auth and reachability +// failures surface as an error from NewStream rather than landing on +// the Errors channel later. +// +// Two goroutines back each Stream: a pull loop and a renew loop. Both +// exit when the context passed to NewStream is cancelled or when Close +// is called. Close is idempotent and bounded — see Stream.Close. +// +// Events is closed exactly when the Stream stops. Ranging over Events +// is safe; a closed channel terminates the loop without a Close call. +// Errors is also closed at stop time. Both channels are buffered (16 +// slots by default); sends to Errors are non-blocking so a stalled +// consumer drops older errors rather than the pull loop blocking on +// log output. +// +// The decoded Event preserves the wire form (Topic, raw Source and +// Data maps) so callers can fall back to inspecting non-standard +// payloads when Kind is KindUnknown. +// +// # Reconnect +// +// On ReconnectAfterFailures consecutive PullMessages failures the +// Stream silently recreates its pull-point subscription. ONVIF cameras +// replay each property's current value with PropertyInitialized on a +// new subscription; Events delivered between recreate and the first +// non-Initialized event carry Event.AfterReconnect=true so consumers +// can suppress duplicate handling. +// +// Set Options.DisableReconnect=true to opt out of recreate; the pull +// loop will retry against the original subscription until ctx cancel. +// +// # Topic classification +// +// Classify maps ONVIF topic strings to a small set of normalized Kind +// values across AXIS, Hikvision, Avigilon, Hanwha, Bosch and Dahua. See +// topics.go for the verified mapping table with public-doc citations. package stream diff --git a/event/stream/stream.go b/event/stream/stream.go index f10c9c5..520843d 100644 --- a/event/stream/stream.go +++ b/event/stream/stream.go @@ -31,55 +31,71 @@ const maxResponseBytes = 10 << 20 // elapses, so a missed unsubscribe is at worst cosmetic. const closeUnsubscribeTimeout = 5 * time.Second -// Options configures a Stream. The zero value is usable; defaultOptions -// fills in production-sensible defaults for any unset field. +// Options configures a Stream. +// +// Zero-value policy: every duration / int field treats zero as "use the +// default". To opt out of reconnect entirely set DisableReconnect=true +// (sentinel `ReconnectAfterFailures=0` would otherwise collide with the +// default-injection policy). To get a synchronous (unbuffered) channel +// pair set BufferSize=-1. type Options struct { - // DeviceID identifies the camera in emitted Events. Recommended so a - // single channel can fan in multiple cameras. Empty is allowed. + // DeviceID identifies the camera in emitted Events. Recommended so + // a single channel can fan in multiple cameras. Empty is allowed. DeviceID string - // TopicFilter is the raw ONVIF ConcreteSet TopicExpression filter - // passed to CreatePullPointSubscription. The empty string means no + // RawTopicFilter is the raw ONVIF ConcreteSet TopicExpression + // filter passed to CreatePullPointSubscription. Empty means no // filter — required for AXIS, accepted by every other vendor we - // support. Callers should normally leave this empty and rely on - // Classify for routing. - TopicFilter string - // PullTimeout is the server-side wait time in each PullMessages call - // (xsd:duration). The camera returns early when messages are - // available; otherwise it returns empty after this timeout. Default: - // 5s. + // support. The name carries 'Raw' because the value is fed verbatim + // into the SOAP envelope: callers should normally leave it empty + // and rely on Classify for routing rather than ask the camera to + // filter server-side, which is fragile across vendors. + RawTopicFilter string + // PullTimeout is the server-side wait time in each PullMessages + // call (xsd:duration). The camera returns early when messages are + // available; otherwise it returns empty after this timeout. Zero + // means default (5s). PullTimeout time.Duration // MessageLimit caps the number of NotificationMessage entries - // returned per PullMessages call. Default: 10. + // returned per PullMessages call. Zero means default (32). A busy + // AXIS with many configured inputs can burst beyond 10 per pull; + // 32 covers that without significantly enlarging quiet pulls. MessageLimit int // InitialTermination is the requested subscription lifetime passed // to CreatePullPointSubscription. The renew loop refreshes well - // before this expires. Default: 60s. + // before this expires. Zero means default (60s). InitialTermination time.Duration // RenewMargin is how long before InitialTermination expiry the // renew loop fires. Larger margins tolerate slower networks at the - // cost of more renew SOAP calls. Default: 10s. + // cost of more renew SOAP calls. Zero means default (10s). RenewMargin time.Duration // ReconnectAfterFailures is the consecutive PullMessages failure // count that triggers a CreatePullPointSubscription recreate. The // camera or pull-point can die for many reasons (camera reboot, // subscription garbage-collected after a renew miss, intermediate // NAT timeout); rebuilding the subscription is the only reliable - // recovery. Default: 3. + // recovery. Zero means default (3). To disable reconnect entirely + // set DisableReconnect=true. ReconnectAfterFailures int + // DisableReconnect skips automatic CreatePullPointSubscription + // recreate. The pull loop will continue retrying against the + // original endpoint until ctx is cancelled. Useful for tests or + // callers managing recovery externally. + DisableReconnect bool // RetryBackoff is the initial sleep between a pull/recreate failure // and the next attempt. Recreate failures double this up to a 30s - // ceiling. Default: 1s. + // ceiling. Zero means default (1s). RetryBackoff time.Duration // BufferSize is the buffer size of the Events and Errors channels. // Larger buffers absorb consumer hiccups at the cost of memory. - // Default: 16. + // Zero means default (16); use -1 for unbuffered (synchronous) + // channels. BufferSize int } func defaultOptions() Options { return Options{ PullTimeout: 5 * time.Second, - MessageLimit: 10, + MessageLimit: 32, InitialTermination: 60 * time.Second, RenewMargin: 10 * time.Second, ReconnectAfterFailures: 3, @@ -108,11 +124,16 @@ func (o Options) withDefaults() Options { if o.RetryBackoff > 0 { d.RetryBackoff = o.RetryBackoff } - if o.BufferSize > 0 { + // BufferSize: zero -> default; negative -> 0 (unbuffered). + switch { + case o.BufferSize > 0: d.BufferSize = o.BufferSize + case o.BufferSize < 0: + d.BufferSize = 0 } d.DeviceID = o.DeviceID - d.TopicFilter = o.TopicFilter + d.RawTopicFilter = o.RawTopicFilter + d.DisableReconnect = o.DisableReconnect return d } @@ -280,7 +301,7 @@ func (s *Stream) pullLoop(ctx context.Context) { if err != nil { s.surfaceError(ErrPullFailed{Err: err}) failures++ - if failures >= s.opts.ReconnectAfterFailures { + if !s.opts.DisableReconnect && failures >= s.opts.ReconnectAfterFailures { justRecreated, cont := s.attemptRecreate(ctx, &failures, &recreateBackoff) if !cont { return @@ -300,7 +321,7 @@ func (s *Stream) pullLoop(ctx context.Context) { recreateBackoff = s.opts.RetryBackoff observedAt := s.now() for _, m := range msgs { - ev := Decode(m, s.opts.DeviceID, observedAt) + ev := decode(m, s.opts.DeviceID, observedAt) if afterReconnect { ev.AfterReconnect = true // ONVIF replays current state with @@ -399,11 +420,11 @@ func sleepCtx(ctx context.Context, d time.Duration) bool { func createPullPoint(c caller, opts Options) (string, error) { term := xsd.String(durationToXSD(opts.InitialTermination)) req := event.CreatePullPointSubscription{InitialTerminationTime: &term} - if opts.TopicFilter != "" { + if opts.RawTopicFilter != "" { req.Filter = &event.FilterType{ TopicExpression: &event.TopicExpressionType{ Dialect: xsd.String("http://www.onvif.org/ver10/tev/topicExpression/ConcreteSet"), - TopicKinds: xsd.String(opts.TopicFilter), + TopicKinds: xsd.String(opts.RawTopicFilter), }, } } diff --git a/event/stream/stream_test.go b/event/stream/stream_test.go index 010f10a..b7343fe 100644 --- a/event/stream/stream_test.go +++ b/event/stream/stream_test.go @@ -326,7 +326,7 @@ func TestStream_PullErrorSurfacedOnErrorsChannel(t *testing.T) { func TestStream_OptionsApplyDefaults(t *testing.T) { o := defaultOptions() assert.Equal(t, 5*time.Second, o.PullTimeout) - assert.Equal(t, 10, o.MessageLimit) + assert.Equal(t, 32, o.MessageLimit) assert.Equal(t, 60*time.Second, o.InitialTermination) assert.Equal(t, 16, o.BufferSize) } diff --git a/examples/event/stream/main.go b/examples/event/stream/main.go index a795ed4..da5cc6e 100644 --- a/examples/event/stream/main.go +++ b/examples/event/stream/main.go @@ -73,7 +73,7 @@ func main() { s, err := stream.NewStream(ctx, dev, stream.Options{ DeviceID: *deviceID, - TopicFilter: *filter, + RawTopicFilter: *filter, PullTimeout: *pullTimeout, }) if err != nil {