diff --git a/cmd/yamcs-recorder/README.md b/cmd/yamcs-recorder/README.md index 77500951..3a7240fa 100644 --- a/cmd/yamcs-recorder/README.md +++ b/cmd/yamcs-recorder/README.md @@ -1,10 +1,12 @@ # yamcs-recorder -Records the TELEMETERED parameters of a YAMCS instance into TimescaleDB, through the yamcs-grpc plugin. +Records the TELEMETERED parameters of a YAMCS instance into TimescaleDB, through the yamcs-grpc plugin. With +`--otlp`, it also exports the instance's events as OpenTelemetry log records to an OpenTelemetry collector, which +forwards them to Loki. ```sh PGPASSWORD=... go run ./cmd/yamcs-recorder --yamcs localhost:8091 --instance fprime-project \ - --postgresql 'postgres://postgres@localhost:5432/hermes?sslmode=disable' + --postgresql 'postgres://postgres@localhost:5432/hermes?sslmode=disable' --otlp localhost:4317 ``` On start it creates the `timescaledb` extension, the `parameters` table and the `parameter_values` hypertable if @@ -16,18 +18,65 @@ fprime-yamcs puts the CCSDS and F Prime packet header fields. To record them any YAMCS sends each parameter's cached value first. After a restart, that value is dropped if it is already stored. -The recorder exits non-zero when the subscription ends, for example when YAMCS stops or its instance restarts, so -run it under a supervisor that restarts it. +The recorder exits non-zero when its parameter or event subscription ends, for example when YAMCS stops or its +instance restarts, so run it under a supervisor that restarts it. If a write to TimescaleDB fails, for example while the database is down, the values in that message are dropped. The recorder keeps running and writes again once the database is back. Each failed write adds one to `insert_errors` on the `recorder stats` line, which it logs every 30 seconds. +## Events + +`--otlp` is the `host:port` of the collector's OTLP gRPC receiver, such as `localhost:4317` for the root +`docker-compose.yml`. The recorder connects without TLS, so it cannot reach a collector that requires TLS. Without +`--otlp` it records parameters only and does not subscribe to events. + +Each event becomes one log record: + +- The service name is the YAMCS instance, and the instrumentation scope is `yamcs-recorder`. +- The timestamp is the event's generation time (when its source raised it) and the observed timestamp its reception + time (when YAMCS received it). When YAMCS sends no reception time, the OpenTelemetry SDK uses the time the recorder + handed it the event. +- The body is the event's message. +- The severity number places the YAMCS severity in an OpenTelemetry band, and Loki gives the record that band's + level: `INFO` 9 (info), `WATCH` 13 and `WARNING` 14 (warn), `DISTRESS` 17 (error), `CRITICAL` 21 and `SEVERE` 22 + (fatal, which Grafana shows as critical). The numbers keep YAMCS's order, so queries can filter on them. + `WARNING_NEW`, the number YAMCS plans to move `WARNING` to, also gets 14, and the deprecated `ERROR` gets 22. There + is no severity text, because Loki would take the level from that text and does not know `WATCH`, `DISTRESS` or + `SEVERE`. +- `yamcs.severity`, `yamcs.source`, `yamcs.type`, `yamcs.seq_number` and `yamcs.created_by` hold YAMCS's own fields. + `yamcs.type` and `yamcs.created_by` are there only when YAMCS sets them. `yamcs.severity` is the severity name, + with `WARNING_NEW` named `WARNING` and `ERROR` named `SEVERE`. +- Each entry of the event's `extra` becomes an attribute `yamcs.extra.`; for F Prime events those are the event + arguments, `fprime_event_id`, `fprime_event_name` and `fprime_severity`. The prefix stops an argument named + `level` or `severity` from replacing the level Loki gives the record. + +Only events raised while the recorder runs are exported; it does not read missed ones from the YAMCS archive. The +recorder sends records in batches about once a second, and on exit sends the rest before it stops. If the collector +is unreachable, the exporter retries each batch for about 10 seconds, then logs an `OpenTelemetry error` line and +drops it. Exit can then take up to about 30 seconds, and the recorder may also log `failed to flush events`. If more +than 2048 events are waiting while the collector is unreachable, the SDK drops the oldest without logging. `events` +on the `recorder stats` line counts the events handed to the SDK, not the ones the collector accepted. + +Loki drops an event whose generation time is more than 7 days old, more than 10 minutes in the future, or more than +an hour older than the newest event it already has for that instance, such as an event downlinked late. The +collector logs `Exporting failed` for it, and the recorder does not see the error. + +In Grafana, query events through the Loki datasource. Loki indexes the service name as the `service_name` label and +keeps the other fields as structured metadata, with dots replaced by underscores. For example, the events of one +F Prime event type, and the errors and worse: + +```logql +{service_name="fprime-project"} | yamcs_type="CdhCore.cmdDisp.OpCodeDispatched" +{service_name="fprime-project"} | severity_number >= 17 +``` + ## Tests The store tests need `HERMES_TEST_TIMESCALE_DSN`, a `postgres://` URL whose user can create databases. The -subscription test needs `YAMCS_GRPC_ADDRESS` and reads `YAMCS_INSTANCE` (default `fprime-project`). Without them, -those tests skip. +subscription tests need `YAMCS_GRPC_ADDRESS` and read `YAMCS_INSTANCE` (default `fprime-project`). Without them, +those tests skip. The event subscription test raises one event through YAMCS's `CreateEvent` call, and the event +stays in that instance's archive. ```sh HERMES_TEST_TIMESCALE_DSN='postgres://postgres:password@localhost:5432/postgres?sslmode=disable' \ diff --git a/cmd/yamcs-recorder/main.go b/cmd/yamcs-recorder/main.go index 54780efb..c4c427c2 100644 --- a/cmd/yamcs-recorder/main.go +++ b/cmd/yamcs-recorder/main.go @@ -1,7 +1,9 @@ // Command yamcs-recorder records the TELEMETERED parameters of a YAMCS // instance into TimescaleDB, by default leaving out the packet header fields -// fprime-yamcs defines. It exits non-zero when the subscription ends, so run it -// under a supervisor that restarts it. +// fprime-yamcs defines. Given --otlp, it also exports the instance's events as +// OpenTelemetry log records to a collector. It exits non-zero when its +// parameter or event subscription ends, so run it under a supervisor that +// restarts it. Events raised while it is not running are not exported. package main import ( @@ -10,10 +12,12 @@ import ( "errors" "fmt" "log/slog" + "net" "net/url" "os" "os/signal" "slices" + "strconv" "strings" "sync/atomic" "syscall" @@ -21,11 +25,19 @@ import ( "github.com/lib/pq" flag "github.com/spf13/pflag" + "go.opentelemetry.io/otel" + "go.opentelemetry.io/otel/exporters/otlp/otlplog/otlploggrpc" + otellog "go.opentelemetry.io/otel/log" + sdklog "go.opentelemetry.io/otel/sdk/log" + "go.opentelemetry.io/otel/sdk/resource" + semconv "go.opentelemetry.io/otel/semconv/v1.4.0" "google.golang.org/grpc" + "google.golang.org/grpc/backoff" "github.com/nasa/hermes/internal/recorder/convert" "github.com/nasa/hermes/internal/recorder/store" "github.com/nasa/hermes/internal/recorder/yamcs" + "github.com/nasa/hermes/internal/yamcspb/protobuf/events" "github.com/nasa/hermes/internal/yamcspb/protobuf/processing" ) @@ -36,6 +48,7 @@ type config struct { instance string processor string postgresql string + otlp string includeTopLevel bool } @@ -45,6 +58,8 @@ func main() { flag.StringVar(&cfg.instance, "instance", "fprime-project", "YAMCS instance to record") flag.StringVar(&cfg.processor, "processor", "realtime", "YAMCS processor to subscribe to") flag.StringVar(&cfg.postgresql, "postgresql", "", "TimescaleDB database URL; give the password in PGPASSWORD") + flag.StringVar(&cfg.otlp, "otlp", "", + "host:port of an OpenTelemetry collector's OTLP gRPC receiver; exports events when set") flag.BoolVar(&cfg.includeTopLevel, "include-top-level", false, "Also record parameters directly in a top-level space system, where fprime-yamcs puts its packet header fields") logLevel := flag.String("log-level", "info", "Log level: debug, info, warn or error") @@ -60,6 +75,29 @@ func main() { logger.Error("--postgresql is required") os.Exit(2) } + if cfg.otlp != "" { + // The OTLP exporter accepts a URL like http://localhost:4317 and only + // fails when it exports, so we check for host:port here. SplitHostPort + // alone accepts http://collector, with port //collector. + _, port, err := net.SplitHostPort(cfg.otlp) + if err == nil { + _, err = strconv.ParseUint(port, 10, 16) + } + if err != nil { + logger.Error("--otlp must be host:port", "err", err) + os.Exit(2) + } + } + // The OpenTelemetry SDK hands errors from its background work, such as failed + // exports, to a global handler that prints them with log.Print. We log them + // through our logger instead so they look like our other lines. Unlike insert + // errors, we log each one rather than count it on the stats line, because the + // exporter sends one batch at a time and retries it for up to 10 seconds before + // dropping it. So while the collector is unreachable, this logs about once every + // 10 seconds. + otel.SetErrorHandler(otel.ErrorHandlerFunc(func(err error) { + logger.Error("OpenTelemetry error", "err", err) + })) ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM) // Stop catching signals once the first arrives, so a second Ctrl-C kills the process. @@ -74,8 +112,11 @@ func main() { os.Exit(1) } -// run records until a signal arrives or the subscription ends. +// run records until a signal arrives or either subscription ends. func run(ctx context.Context, cfg config, logger *slog.Logger) error { + ctx, cancel := context.WithCancel(ctx) + defer cancel() + db, err := openDatabase(ctx, cfg.postgresql, logger) if err != nil { return err @@ -95,6 +136,31 @@ func run(ctx context.Context, cfg config, logger *slog.Logger) error { if err != nil { return err } + var eventStream events.EventsApi_SubscribeEventsClient + var eventLogger otellog.Logger + if cfg.otlp == "" { + logger.Info("not exporting events, since --otlp is not set") + } else { + provider, err := newLoggerProvider(ctx, cfg.otlp, cfg.instance) + if err != nil { + return err + } + // On return we shut the provider down, which exports the events it still + // holds. By then ctx has usually been cancelled, so we pass a fresh context. + // We give it no deadline, since the exporter gives up on each send after 10 + // seconds. Even with the collector down, exit takes at most about 30 seconds. + defer func() { + if err := provider.Shutdown(context.Background()); err != nil { + logger.Error("failed to flush events", "err", err) + } + }() + eventLogger = provider.Logger("yamcs-recorder") + eventStream, err = yamcs.SubscribeEvents(ctx, conn, cfg.instance) + if err != nil { + return err + } + logger.Info("exporting events", "otlp", cfg.otlp) + } r := &recorder{ db: db, @@ -105,14 +171,23 @@ func run(ctx context.Context, cfg config, logger *slog.Logger) error { } go r.stats.logEvery(ctx, logger) defer r.stats.log(logger) - for ctx.Err() == nil { - data, err := stream.Recv() - if err != nil { - return fmt.Errorf("parameter subscription ended: %w", err) - } - r.record(ctx, data) + + // Both subscriptions run on ctx, so when the first one ends we cancel ctx to + // end the other. That is also how the event stream ends when the instance + // restarts, since YAMCS leaves it open but silent. We return the first one's + // error, but only after both loops are done, so the deferred flush, stats + // line and database close come after the last emit and insert. + done := make(chan error, 2) + go func() { done <- recordParameters(ctx, stream, r) }() + if eventStream != nil { + go func() { done <- exportEvents(ctx, eventStream, eventLogger, &r.stats) }() } - return ctx.Err() + err = <-done + cancel() + if eventStream != nil { + <-done + } + return err } // openDatabase connects to the database at dbURL and creates or updates the @@ -174,6 +249,58 @@ func inTopLevelSpaceSystem(name string) bool { return strings.Count(name, "/") == 2 } +// newLoggerProvider returns a provider that batches log records and exports +// them over OTLP gRPC to the collector at addr. Its resource, the attributes it +// sends with every batch, names the instance as the service, which Loki turns +// into the service_name label. +func newLoggerProvider(ctx context.Context, addr, instance string) (*sdklog.LoggerProvider, error) { + // gRPC waits up to 2 minutes between reconnects, so after a long collector + // outage we would keep dropping events that long once it is back. We lower + // the limit to 5 seconds. ConnectParams also replaces gRPC's 20 second + // connect timeout, so we pass that default again. + reconnect := backoff.DefaultConfig + reconnect.MaxDelay = 5 * time.Second + exporter, err := otlploggrpc.New(ctx, + otlploggrpc.WithEndpoint(addr), + // The collector in docker-compose.yml doesn't use TLS. + otlploggrpc.WithInsecure(), + otlploggrpc.WithDialOption(grpc.WithConnectParams(grpc.ConnectParams{Backoff: reconnect, MinConnectTimeout: 20 * time.Second})), + ) + if err != nil { + return nil, fmt.Errorf("failed to create OTLP exporter: %w", err) + } + return sdklog.NewLoggerProvider( + sdklog.WithResource(resource.NewSchemaless(semconv.ServiceNameKey.String(instance))), + sdklog.WithProcessor(sdklog.NewBatchProcessor(exporter)), + ), nil +} + +// recordParameters stores the values that stream delivers until it ends. +func recordParameters(ctx context.Context, stream processing.ProcessingApi_SubscribeParametersClient, r *recorder) error { + for ctx.Err() == nil { + data, err := stream.Recv() + if err != nil { + return fmt.Errorf("parameter subscription ended: %w", err) + } + r.record(ctx, data) + } + return ctx.Err() +} + +// exportEvents emits each event that stream delivers as a log record until the +// stream ends. The provider behind eventLogger batches and exports them. +func exportEvents(ctx context.Context, stream events.EventsApi_SubscribeEventsClient, eventLogger otellog.Logger, st *stats) error { + for ctx.Err() == nil { + e, err := stream.Recv() + if err != nil { + return fmt.Errorf("event subscription ended: %w", err) + } + eventLogger.Emit(ctx, convert.LogRecord(e)) + st.events.Add(1) + } + return ctx.Err() +} + // recorder stores the values that arrive on one parameter subscription. type recorder struct { db *sql.DB @@ -242,6 +369,7 @@ func dropStored(rows []store.Row, stored map[store.Parameter]time.Time) (kept [] type stats struct { values atomic.Int64 // parameter values YAMCS sent rows atomic.Int64 // rows inserted + events atomic.Int64 // events handed to the OpenTelemetry SDK unmapped atomic.Int64 // values dropped because YAMCS never sent their parameter's name incomplete atomic.Int64 // values and members dropped for having no generation time or no value alreadyStored atomic.Int64 // rows dropped by dropStored @@ -252,7 +380,7 @@ type stats struct { func (s *stats) log(logger *slog.Logger) { attrs := []any{ - "values", s.values.Load(), "rows", s.rows.Load(), "unmapped", s.unmapped.Load(), + "values", s.values.Load(), "rows", s.rows.Load(), "events", s.events.Load(), "unmapped", s.unmapped.Load(), "incomplete", s.incomplete.Load(), "already_stored", s.alreadyStored.Load(), "insert_errors", s.insertErrors.Load(), } // Swap clears the error, so each line shows only the newest one since the previous line. diff --git a/cmd/yamcs-recorder/main_test.go b/cmd/yamcs-recorder/main_test.go index 77378f12..9582f5a9 100644 --- a/cmd/yamcs-recorder/main_test.go +++ b/cmd/yamcs-recorder/main_test.go @@ -67,9 +67,10 @@ func TestDropStoredKeepsNewerRows(t *testing.T) { func TestStatsLineNamesEachCount(t *testing.T) { var st stats + st.events.Add(4) st.incomplete.Add(3) st.alreadyStored.Add(2) var logs bytes.Buffer st.log(slog.New(slog.NewTextHandler(&logs, nil))) - assert.Contains(t, logs.String(), "unmapped=0 incomplete=3 already_stored=2 insert_errors=0") + assert.Contains(t, logs.String(), "rows=0 events=4 unmapped=0 incomplete=3 already_stored=2 insert_errors=0") } diff --git a/docker-compose.yml b/docker-compose.yml index 95680fb8..443df148 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -20,9 +20,31 @@ services: depends_on: timescaledb: condition: service_started + loki: + condition: service_started volumes: - grafana-data:/var/lib/grafana - ./config:/etc/grafana/provisioning:ro + loki: + image: grafana/loki:latest + restart: unless-stopped + ports: + - "3100:3100" + volumes: + - loki-data:/loki + - ./config/loki.yml:/etc/loki/local-config.yaml:ro + otel-collector: + image: otel/opentelemetry-collector:latest + restart: unless-stopped + ports: + - "4317:4317" + - "4318:4318" + depends_on: + loki: + condition: service_started + volumes: + - ./config/otel-collector.yml:/etc/otelcol/config.yaml:ro volumes: timescale-data: grafana-data: + loki-data: diff --git a/grafana-datasource-plugin/src/README.md b/grafana-datasource-plugin/src/README.md index f9cd77cd..bbd5d038 100644 --- a/grafana-datasource-plugin/src/README.md +++ b/grafana-datasource-plugin/src/README.md @@ -74,7 +74,7 @@ Notes: #### Builder: Events -The Events query type is hidden because yamcs-recorder only records telemetry. +The Events query type is hidden because yamcs-recorder records only telemetry in TimescaleDB. With `--otlp`, it exports events through an OpenTelemetry collector to Loki, so query them through Grafana's Loki datasource.
diff --git a/internal/recorder/convert/convert.go b/internal/recorder/convert/convert.go index 07ee5ad7..abfba861 100644 --- a/internal/recorder/convert/convert.go +++ b/internal/recorder/convert/convert.go @@ -1,5 +1,5 @@ // Package convert turns YAMCS parameter subscription messages into recorder -// store rows. +// store rows, and YAMCS events into OpenTelemetry log records. package convert import ( diff --git a/internal/recorder/convert/events.go b/internal/recorder/convert/events.go new file mode 100644 index 00000000..5044b3cb --- /dev/null +++ b/internal/recorder/convert/events.go @@ -0,0 +1,88 @@ +package convert + +import ( + "maps" + "slices" + + otellog "go.opentelemetry.io/otel/log" + + "github.com/nasa/hermes/internal/yamcspb/protobuf/events" +) + +// severityNumbers puts each YAMCS severity in an OpenTelemetry band (info, +// warn, error or fatal), which Loki turns into the level Grafana shows, and +// numbers them in YAMCS's order so queries can filter on severity_number. We +// put CRITICAL in fatal so it shows apart from DISTRESS, since fprime-yamcs +// sends F Prime's FATAL as CRITICAL and WARNING_HI as DISTRESS. YAMCS already +// sends its deprecated ERROR as SEVERE, so we number ERROR the same. +var severityNumbers = map[events.Event_EventSeverity]otellog.Severity{ + events.Event_INFO: otellog.SeverityInfo1, + events.Event_WATCH: otellog.SeverityWarn1, + events.Event_WARNING: otellog.SeverityWarn2, + events.Event_WARNING_NEW: otellog.SeverityWarn2, + events.Event_DISTRESS: otellog.SeverityError1, + events.Event_CRITICAL: otellog.SeverityFatal1, + events.Event_SEVERE: otellog.SeverityFatal2, + events.Event_ERROR: otellog.SeverityFatal2, +} + +// LogRecord converts an event to an OpenTelemetry log record with the message +// as its body. The timestamp is the generation time and the observed timestamp +// the reception time. The SDK sets a missing observed timestamp to the time we +// emit the record, and Loki uses the observed timestamp when the timestamp is +// missing, so we keep an event that lacks either time and leave that timestamp +// unset. The other fields become attributes under "yamcs.", and each extra +// entry becomes one under "yamcs.extra.". +func LogRecord(e *events.Event) otellog.Record { + var r otellog.Record + if e.GenerationTime != nil { + r.SetTimestamp(e.GetGenerationTime().AsTime()) + } + if e.ReceptionTime != nil { + r.SetObservedTimestamp(e.GetReceptionTime().AsTime()) + } + // Loki reads the level from the severity text when there is one, and it does + // not know WATCH, DISTRESS or SEVERE, so Grafana would show those with no + // level. We leave the text unset and keep the YAMCS name in yamcs.severity. + r.SetSeverity(severityNumbers[e.GetSeverity()]) + r.SetBody(otellog.StringValue(e.GetMessage())) + + r.AddAttributes( + otellog.String("yamcs.severity", severityName(e.GetSeverity())), + otellog.String("yamcs.source", e.GetSource()), + otellog.Int64("yamcs.seq_number", int64(e.GetSeqNumber())), + ) + // Not every YAMCS event has a type, and only events posted through the YAMCS + // API have createdBy, so we add yamcs.type and yamcs.created_by only when set. + if e.GetType() != "" { + r.AddAttributes(otellog.String("yamcs.type", e.GetType())) + } + if e.GetCreatedBy() != "" { + r.AddAttributes(otellog.String("yamcs.created_by", e.GetCreatedBy())) + } + // fprime-yamcs puts each F Prime event argument in extra under its own name. + // Loki reads the level from an attribute named level or severity, among + // others, before the severity number, and the SDK keeps only the last of two + // attributes with the same name. So we put extra under "yamcs.extra.", where + // an argument can neither set the level nor replace a YAMCS field. We add the + // entries in key order so the same event always gives an identical record, + // which Loki stores only once if two recorders export it. + for _, k := range slices.Sorted(maps.Keys(e.GetExtra())) { + r.AddAttributes(otellog.String("yamcs.extra."+k, e.GetExtra()[k])) + } + return r +} + +// severityName returns the name of s, with WARNING_NEW named WARNING and ERROR +// named SEVERE. YAMCS already sends ERROR as SEVERE. It plans to move WARNING +// to number 4, which our copy of events.proto calls WARNING_NEW, so its +// warnings will then arrive as WARNING_NEW. +func severityName(s events.Event_EventSeverity) string { + switch s { + case events.Event_WARNING_NEW: + return events.Event_WARNING.String() + case events.Event_ERROR: + return events.Event_SEVERE.String() + } + return s.String() +} diff --git a/internal/recorder/convert/events_test.go b/internal/recorder/convert/events_test.go new file mode 100644 index 00000000..3930a310 --- /dev/null +++ b/internal/recorder/convert/events_test.go @@ -0,0 +1,137 @@ +package convert + +import ( + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + otellog "go.opentelemetry.io/otel/log" + "google.golang.org/protobuf/encoding/protojson" + + "github.com/nasa/hermes/internal/yamcspb/protobuf/events" +) + +func parseEvent(t *testing.T, s string) *events.Event { + t.Helper() + e := &events.Event{} + require.NoError(t, protojson.Unmarshal([]byte(s), e)) + return e +} + +// attributes returns the attributes of r as strings and int64s, failing on +// duplicate keys and on other kinds. +func attributes(t *testing.T, r otellog.Record) map[string]any { + t.Helper() + out := make(map[string]any) + r.WalkAttributes(func(kv otellog.KeyValue) bool { + require.NotContains(t, out, kv.Key) + switch kv.Value.Kind() { + case otellog.KindString: + out[kv.Key] = kv.Value.AsString() + case otellog.KindInt64: + out[kv.Key] = kv.Value.AsInt64() + default: + t.Errorf("attribute %s has kind %s", kv.Key, kv.Value.Kind()) + } + return true + }) + return out +} + +// dispatchedEvent comes from the YAMCS archive of the fprime-project instance. +// fprime-yamcs's event processor decoded it from an F Prime event and posted it +// through the YAMCS API, which set createdBy. +const dispatchedEvent = `{"source": "FPrimeEventProcessor", "generationTime": "2026-10-05T22:56:34.620Z", + "receptionTime": "2026-10-05T22:56:34.591Z", "seqNumber": 6, "type": "CdhCore.cmdDisp.OpCodeDispatched", + "message": "[OpCodeDispatched] Opcode 0x1000000 dispatched to port 4", "severity": "INFO", "createdBy": "guest", + "extra": {"Opcode": "16777216", "fprime_event_id": "16777217", "fprime_event_name": "CdhCore.cmdDisp.OpCodeDispatched", + "fprime_severity": "COMMAND", "port": "4"}}` + +// seqCountJumpEvent comes from the same archive. fprime-yamcs's packet +// preprocessor raises it inside YAMCS, so it has no extra and no createdBy. +const seqCountJumpEvent = `{"source": "FprimePacketPreprocessor", "generationTime": "2026-09-29T22:54:28.095Z", + "receptionTime": "2026-09-29T22:54:28.095Z", "seqNumber": 1, "type": "SEQ_COUNT_JUMP", + "message": "Sequence count jump for APID: 4 old seq: 0 newseq: 0", "severity": "WARNING"}` + +func TestLogRecordKeepsYamcsFields(t *testing.T) { + r := LogRecord(parseEvent(t, dispatchedEvent)) + assert.Equal(t, utc("2026-10-05T22:56:34.620Z"), r.Timestamp()) + assert.Equal(t, utc("2026-10-05T22:56:34.591Z"), r.ObservedTimestamp()) + assert.Equal(t, otellog.SeverityInfo1, r.Severity()) + assert.Empty(t, r.SeverityText()) + assert.Equal(t, "[OpCodeDispatched] Opcode 0x1000000 dispatched to port 4", r.Body().AsString()) + assert.Equal(t, map[string]any{ + "yamcs.severity": "INFO", + "yamcs.source": "FPrimeEventProcessor", + "yamcs.type": "CdhCore.cmdDisp.OpCodeDispatched", + "yamcs.seq_number": int64(6), + "yamcs.created_by": "guest", + "yamcs.extra.Opcode": "16777216", + "yamcs.extra.fprime_event_id": "16777217", + "yamcs.extra.fprime_event_name": "CdhCore.cmdDisp.OpCodeDispatched", + "yamcs.extra.fprime_severity": "COMMAND", + "yamcs.extra.port": "4", + }, attributes(t, r)) +} + +func TestLogRecordWithoutTypeExtraOrCreatedBy(t *testing.T) { + r := LogRecord(parseEvent(t, seqCountJumpEvent)) + assert.Equal(t, otellog.SeverityWarn2, r.Severity()) + assert.Equal(t, map[string]any{ + "yamcs.severity": "WARNING", + "yamcs.source": "FprimePacketPreprocessor", + "yamcs.type": "SEQ_COUNT_JUMP", + "yamcs.seq_number": int64(1), + }, attributes(t, r)) + + // We clear the archived event's type to check that yamcs.type is left out. + e := parseEvent(t, seqCountJumpEvent) + e.Type = nil + assert.NotContains(t, attributes(t, LogRecord(e)), "yamcs.type") +} + +// We write out the severity numbers rather than otellog's names for them, +// since Loki stores the numbers as severity_number and queries filter on them. +func TestLogRecordSeverities(t *testing.T) { + for s, want := range map[events.Event_EventSeverity]struct { + number otellog.Severity + name string + }{ + events.Event_INFO: {9, "INFO"}, + events.Event_WATCH: {13, "WATCH"}, + events.Event_WARNING: {14, "WARNING"}, + events.Event_WARNING_NEW: {14, "WARNING"}, + events.Event_DISTRESS: {17, "DISTRESS"}, + events.Event_CRITICAL: {21, "CRITICAL"}, + events.Event_SEVERE: {22, "SEVERE"}, + events.Event_ERROR: {22, "SEVERE"}, + } { + e := parseEvent(t, seqCountJumpEvent) + e.Severity = s.Enum() + r := LogRecord(e) + assert.Equal(t, want.number, r.Severity(), s.String()) + assert.Equal(t, want.name, attributes(t, r)["yamcs.severity"], s.String()) + } + assert.Len(t, events.Event_EventSeverity_name, 8, "a new YAMCS severity needs a number in severityNumbers") + + // An unset severity reads as INFO, the default events.proto declares. + e := parseEvent(t, seqCountJumpEvent) + e.Severity = nil + r := LogRecord(e) + assert.Equal(t, otellog.Severity(9), r.Severity()) + assert.Equal(t, "INFO", attributes(t, r)["yamcs.severity"]) +} + +// We keep events that lack a time, leaving that timestamp unset. +func TestLogRecordTimes(t *testing.T) { + e := parseEvent(t, dispatchedEvent) + e.ReceptionTime = nil + r := LogRecord(e) + assert.Equal(t, utc("2026-10-05T22:56:34.620Z"), r.Timestamp()) + assert.Zero(t, r.ObservedTimestamp()) + + e.GenerationTime = nil + r = LogRecord(e) + assert.Zero(t, r.Timestamp()) + assert.Equal(t, "[OpCodeDispatched] Opcode 0x1000000 dispatched to port 4", r.Body().AsString()) +} diff --git a/internal/recorder/yamcs/yamcs.go b/internal/recorder/yamcs/yamcs.go index adaa8e29..be252ec6 100644 --- a/internal/recorder/yamcs/yamcs.go +++ b/internal/recorder/yamcs/yamcs.go @@ -1,6 +1,6 @@ // Package yamcs is the recorder's client for YAMCS, which it reaches through // the yamcs-grpc plugin. It lists the parameters to record and streams their -// values. +// values and the instance's events. package yamcs import ( @@ -16,6 +16,7 @@ import ( "google.golang.org/protobuf/proto" "github.com/nasa/hermes/internal/yamcspb/protobuf" + "github.com/nasa/hermes/internal/yamcspb/protobuf/events" "github.com/nasa/hermes/internal/yamcspb/protobuf/mdb" "github.com/nasa/hermes/internal/yamcspb/protobuf/processing" "github.com/nasa/hermes/internal/yamcspb/protobuf/services" @@ -160,3 +161,25 @@ func (s subscription) Recv() (*processing.SubscribeParametersData, error) { } return data, err } + +// SubscribeEvents streams instance's events as YAMCS raises them, without +// replaying its archive. When the instance restarts, YAMCS leaves the event +// stream open but silent, as it does the parameter stream. Subscribe already +// watches the processor, so we don't watch it again here, and the caller +// should end the event stream when the parameter stream ends. +func SubscribeEvents(ctx context.Context, conn grpc.ClientConnInterface, instance string) (events.EventsApi_SubscribeEventsClient, error) { + stream, err := events.NewEventsApiClient(conn).SubscribeEvents(ctx) + if err != nil { + return nil, fmt.Errorf("failed to open event subscription: %w", err) + } + err = stream.Send(&events.SubscribeEventsRequest{Instance: proto.String(instance)}) + // Send returns io.EOF when YAMCS has already ended the stream, and only + // Recv returns YAMCS's error, so we ask Recv for it. + if errors.Is(err, io.EOF) { + _, err = stream.Recv() + } + if err != nil { + return nil, fmt.Errorf("failed to send event subscription: %w", err) + } + return stream, nil +} diff --git a/internal/recorder/yamcs/yamcs_test.go b/internal/recorder/yamcs/yamcs_test.go index cebe2875..d08268e2 100644 --- a/internal/recorder/yamcs/yamcs_test.go +++ b/internal/recorder/yamcs/yamcs_test.go @@ -10,6 +10,9 @@ import ( "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" + "google.golang.org/protobuf/proto" + + "github.com/nasa/hermes/internal/yamcspb/protobuf/events" ) // Runs against a YAMCS server with the yamcs-grpc plugin and live telemetry, @@ -52,3 +55,39 @@ func TestSubscribeTelemetered(t *testing.T) { } assert.Len(t, mapped, len(names)) } + +// Raises one event with CreateEvent. It stays in the instance's archive. +func TestSubscribeEvents(t *testing.T) { + addr := os.Getenv("YAMCS_GRPC_ADDRESS") + if addr == "" { + t.Skip("YAMCS_GRPC_ADDRESS not set") + } + conn, err := Dial(addr) + require.NoError(t, err) + defer conn.Close() + ctx, cancel := context.WithTimeout(context.Background(), 20*time.Second) + defer cancel() + instance := cmp.Or(os.Getenv("YAMCS_INSTANCE"), "fprime-project") + + stream, err := SubscribeEvents(ctx, conn, instance) + require.NoError(t, err) + // YAMCS does not acknowledge the subscription, so give it time to start. + time.Sleep(time.Second) + created, err := events.NewEventsApiClient(conn).CreateEvent(ctx, &events.CreateEventRequest{ + Instance: proto.String(instance), + Source: proto.String("yamcs-recorder test"), + Message: proto.String(t.Name()), + }) + require.NoError(t, err) + + // Other sources may raise events first. + for { + e, err := stream.Recv() + require.NoError(t, err) + if e.GetSource() == created.GetSource() && e.GetSeqNumber() == created.GetSeqNumber() { + t.Logf("received %s #%d: %s", e.GetSource(), e.GetSeqNumber(), e.GetMessage()) + assert.Equal(t, t.Name(), e.GetMessage()) + return + } + } +} diff --git a/internal/yamcspb/protobuf/events/events.pb.go b/internal/yamcspb/protobuf/events/events.pb.go new file mode 100644 index 00000000..a1920106 --- /dev/null +++ b/internal/yamcspb/protobuf/events/events.pb.go @@ -0,0 +1,324 @@ +// Code generated by protoc-gen-go. DO NOT EDIT. +// versions: +// protoc-gen-go v1.36.11 +// protoc v7.36.0 +// source: yamcs/protobuf/events/events.proto + +package events + +import ( + protoreflect "google.golang.org/protobuf/reflect/protoreflect" + protoimpl "google.golang.org/protobuf/runtime/protoimpl" + timestamppb "google.golang.org/protobuf/types/known/timestamppb" + reflect "reflect" + sync "sync" + unsafe "unsafe" +) + +const ( + // Verify that this generated code is sufficiently up-to-date. + _ = protoimpl.EnforceVersion(20 - protoimpl.MinVersion) + // Verify that runtime/protoimpl is sufficiently up-to-date. + _ = protoimpl.EnforceVersion(protoimpl.MaxVersion - 20) +) + +// The severity levels, in order are: +// INFO, WATCH, WARNING, DISTRESS, CRITICAL, SEVERE. +// +// A migration is underway to fully move away from the legacy +// INFO, WARNING, ERROR levels. +type Event_EventSeverity int32 + +const ( + Event_INFO Event_EventSeverity = 0 + Event_WARNING Event_EventSeverity = 1 + // Legacy, avoid use. + // + // Deprecated: Marked as deprecated in yamcs/protobuf/events/events.proto. + Event_ERROR Event_EventSeverity = 2 + Event_WATCH Event_EventSeverity = 3 + // Placeholder for future WARNING constant. + // (correctly sorted between WATCH and DISTRESS) + // + // Most clients can ignore, this state is here + // to give Protobuf clients (Python Client, Yamcs Studio) + // the time to add a migration for supporting both WARNING + // and WARNING_NEW (Protobuf serializes the number). + // + // Then in a later phase, we move from: + // WARNING=1, WARNING_NEW=4 + // + // To: + // WARNING_OLD=1, WARNING=4 + // + // (which is a transparent change to JSON clients) + Event_WARNING_NEW Event_EventSeverity = 4 + Event_DISTRESS Event_EventSeverity = 5 + Event_CRITICAL Event_EventSeverity = 6 + Event_SEVERE Event_EventSeverity = 7 +) + +// Enum value maps for Event_EventSeverity. +var ( + Event_EventSeverity_name = map[int32]string{ + 0: "INFO", + 1: "WARNING", + 2: "ERROR", + 3: "WATCH", + 4: "WARNING_NEW", + 5: "DISTRESS", + 6: "CRITICAL", + 7: "SEVERE", + } + Event_EventSeverity_value = map[string]int32{ + "INFO": 0, + "WARNING": 1, + "ERROR": 2, + "WATCH": 3, + "WARNING_NEW": 4, + "DISTRESS": 5, + "CRITICAL": 6, + "SEVERE": 7, + } +) + +func (x Event_EventSeverity) Enum() *Event_EventSeverity { + p := new(Event_EventSeverity) + *p = x + return p +} + +func (x Event_EventSeverity) String() string { + return protoimpl.X.EnumStringOf(x.Descriptor(), protoreflect.EnumNumber(x)) +} + +func (Event_EventSeverity) Descriptor() protoreflect.EnumDescriptor { + return file_yamcs_protobuf_events_events_proto_enumTypes[0].Descriptor() +} + +func (Event_EventSeverity) Type() protoreflect.EnumType { + return &file_yamcs_protobuf_events_events_proto_enumTypes[0] +} + +func (x Event_EventSeverity) Number() protoreflect.EnumNumber { + return protoreflect.EnumNumber(x) +} + +// Deprecated: Do not use. +func (x *Event_EventSeverity) UnmarshalJSON(b []byte) error { + num, err := protoimpl.X.UnmarshalJSONEnum(x.Descriptor(), b) + if err != nil { + return err + } + *x = Event_EventSeverity(num) + return nil +} + +// Deprecated: Use Event_EventSeverity.Descriptor instead. +func (Event_EventSeverity) EnumDescriptor() ([]byte, []int) { + return file_yamcs_protobuf_events_events_proto_rawDescGZIP(), []int{0, 0} +} + +type Event struct { + state protoimpl.MessageState `protogen:"open.v1"` + Source *string `protobuf:"bytes,1,opt,name=source" json:"source,omitempty"` + GenerationTime *timestamppb.Timestamp `protobuf:"bytes,2,opt,name=generationTime" json:"generationTime,omitempty"` + ReceptionTime *timestamppb.Timestamp `protobuf:"bytes,3,opt,name=receptionTime" json:"receptionTime,omitempty"` + SeqNumber *int32 `protobuf:"varint,4,opt,name=seqNumber" json:"seqNumber,omitempty"` + Type *string `protobuf:"bytes,5,opt,name=type" json:"type,omitempty"` + Message *string `protobuf:"bytes,6,opt,name=message" json:"message,omitempty"` + Severity *Event_EventSeverity `protobuf:"varint,7,opt,name=severity,enum=yamcs.protobuf.events.Event_EventSeverity,def=0" json:"severity,omitempty"` + // Set by API when event was posted by a user + CreatedBy *string `protobuf:"bytes,10,opt,name=createdBy" json:"createdBy,omitempty"` + // Additional properties + Extra map[string]string `protobuf:"bytes,11,rep,name=extra" json:"extra,omitempty" protobuf_key:"bytes,1,opt,name=key" protobuf_val:"bytes,2,opt,name=value"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +// Default values for Event fields. +const ( + Default_Event_Severity = Event_INFO +) + +func (x *Event) Reset() { + *x = Event{} + mi := &file_yamcs_protobuf_events_events_proto_msgTypes[0] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *Event) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*Event) ProtoMessage() {} + +func (x *Event) ProtoReflect() protoreflect.Message { + mi := &file_yamcs_protobuf_events_events_proto_msgTypes[0] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use Event.ProtoReflect.Descriptor instead. +func (*Event) Descriptor() ([]byte, []int) { + return file_yamcs_protobuf_events_events_proto_rawDescGZIP(), []int{0} +} + +func (x *Event) GetSource() string { + if x != nil && x.Source != nil { + return *x.Source + } + return "" +} + +func (x *Event) GetGenerationTime() *timestamppb.Timestamp { + if x != nil { + return x.GenerationTime + } + return nil +} + +func (x *Event) GetReceptionTime() *timestamppb.Timestamp { + if x != nil { + return x.ReceptionTime + } + return nil +} + +func (x *Event) GetSeqNumber() int32 { + if x != nil && x.SeqNumber != nil { + return *x.SeqNumber + } + return 0 +} + +func (x *Event) GetType() string { + if x != nil && x.Type != nil { + return *x.Type + } + return "" +} + +func (x *Event) GetMessage() string { + if x != nil && x.Message != nil { + return *x.Message + } + return "" +} + +func (x *Event) GetSeverity() Event_EventSeverity { + if x != nil && x.Severity != nil { + return *x.Severity + } + return Default_Event_Severity +} + +func (x *Event) GetCreatedBy() string { + if x != nil && x.CreatedBy != nil { + return *x.CreatedBy + } + return "" +} + +func (x *Event) GetExtra() map[string]string { + if x != nil { + return x.Extra + } + return nil +} + +var File_yamcs_protobuf_events_events_proto protoreflect.FileDescriptor + +const file_yamcs_protobuf_events_events_proto_rawDesc = "" + + "\n" + + "\"yamcs/protobuf/events/events.proto\x12\x15yamcs.protobuf.events\x1a\x1fgoogle/protobuf/timestamp.proto\"\xd1\x04\n" + + "\x05Event\x12\x16\n" + + "\x06source\x18\x01 \x01(\tR\x06source\x12B\n" + + "\x0egenerationTime\x18\x02 \x01(\v2\x1a.google.protobuf.TimestampR\x0egenerationTime\x12@\n" + + "\rreceptionTime\x18\x03 \x01(\v2\x1a.google.protobuf.TimestampR\rreceptionTime\x12\x1c\n" + + "\tseqNumber\x18\x04 \x01(\x05R\tseqNumber\x12\x12\n" + + "\x04type\x18\x05 \x01(\tR\x04type\x12\x18\n" + + "\amessage\x18\x06 \x01(\tR\amessage\x12L\n" + + "\bseverity\x18\a \x01(\x0e2*.yamcs.protobuf.events.Event.EventSeverity:\x04INFOR\bseverity\x12\x1c\n" + + "\tcreatedBy\x18\n" + + " \x01(\tR\tcreatedBy\x12=\n" + + "\x05extra\x18\v \x03(\v2'.yamcs.protobuf.events.Event.ExtraEntryR\x05extra\x1a8\n" + + "\n" + + "ExtraEntry\x12\x10\n" + + "\x03key\x18\x01 \x01(\tR\x03key\x12\x14\n" + + "\x05value\x18\x02 \x01(\tR\x05value:\x028\x01\"y\n" + + "\rEventSeverity\x12\b\n" + + "\x04INFO\x10\x00\x12\v\n" + + "\aWARNING\x10\x01\x12\r\n" + + "\x05ERROR\x10\x02\x1a\x02\b\x01\x12\t\n" + + "\x05WATCH\x10\x03\x12\x0f\n" + + "\vWARNING_NEW\x10\x04\x12\f\n" + + "\bDISTRESS\x10\x05\x12\f\n" + + "\bCRITICAL\x10\x06\x12\n" + + "\n" + + "\x06SEVERE\x10\aB#\n" + + "\x12org.yamcs.protobufB\vEventsProtoP\x01" + +var ( + file_yamcs_protobuf_events_events_proto_rawDescOnce sync.Once + file_yamcs_protobuf_events_events_proto_rawDescData []byte +) + +func file_yamcs_protobuf_events_events_proto_rawDescGZIP() []byte { + file_yamcs_protobuf_events_events_proto_rawDescOnce.Do(func() { + file_yamcs_protobuf_events_events_proto_rawDescData = protoimpl.X.CompressGZIP(unsafe.Slice(unsafe.StringData(file_yamcs_protobuf_events_events_proto_rawDesc), len(file_yamcs_protobuf_events_events_proto_rawDesc))) + }) + return file_yamcs_protobuf_events_events_proto_rawDescData +} + +var file_yamcs_protobuf_events_events_proto_enumTypes = make([]protoimpl.EnumInfo, 1) +var file_yamcs_protobuf_events_events_proto_msgTypes = make([]protoimpl.MessageInfo, 2) +var file_yamcs_protobuf_events_events_proto_goTypes = []any{ + (Event_EventSeverity)(0), // 0: yamcs.protobuf.events.Event.EventSeverity + (*Event)(nil), // 1: yamcs.protobuf.events.Event + nil, // 2: yamcs.protobuf.events.Event.ExtraEntry + (*timestamppb.Timestamp)(nil), // 3: google.protobuf.Timestamp +} +var file_yamcs_protobuf_events_events_proto_depIdxs = []int32{ + 3, // 0: yamcs.protobuf.events.Event.generationTime:type_name -> google.protobuf.Timestamp + 3, // 1: yamcs.protobuf.events.Event.receptionTime:type_name -> google.protobuf.Timestamp + 0, // 2: yamcs.protobuf.events.Event.severity:type_name -> yamcs.protobuf.events.Event.EventSeverity + 2, // 3: yamcs.protobuf.events.Event.extra:type_name -> yamcs.protobuf.events.Event.ExtraEntry + 4, // [4:4] is the sub-list for method output_type + 4, // [4:4] is the sub-list for method input_type + 4, // [4:4] is the sub-list for extension type_name + 4, // [4:4] is the sub-list for extension extendee + 0, // [0:4] is the sub-list for field type_name +} + +func init() { file_yamcs_protobuf_events_events_proto_init() } +func file_yamcs_protobuf_events_events_proto_init() { + if File_yamcs_protobuf_events_events_proto != nil { + return + } + type x struct{} + out := protoimpl.TypeBuilder{ + File: protoimpl.DescBuilder{ + GoPackagePath: reflect.TypeOf(x{}).PkgPath(), + RawDescriptor: unsafe.Slice(unsafe.StringData(file_yamcs_protobuf_events_events_proto_rawDesc), len(file_yamcs_protobuf_events_events_proto_rawDesc)), + NumEnums: 1, + NumMessages: 2, + NumExtensions: 0, + NumServices: 0, + }, + GoTypes: file_yamcs_protobuf_events_events_proto_goTypes, + DependencyIndexes: file_yamcs_protobuf_events_events_proto_depIdxs, + EnumInfos: file_yamcs_protobuf_events_events_proto_enumTypes, + MessageInfos: file_yamcs_protobuf_events_events_proto_msgTypes, + }.Build() + File_yamcs_protobuf_events_events_proto = out.File + file_yamcs_protobuf_events_events_proto_goTypes = nil + file_yamcs_protobuf_events_events_proto_depIdxs = nil +} diff --git a/internal/yamcspb/protobuf/events/events_service.pb.go b/internal/yamcspb/protobuf/events/events_service.pb.go new file mode 100644 index 00000000..07393d14 --- /dev/null +++ b/internal/yamcspb/protobuf/events/events_service.pb.go @@ -0,0 +1,934 @@ +// Code generated by protoc-gen-go. DO NOT EDIT. +// versions: +// protoc-gen-go v1.36.11 +// protoc v7.36.0 +// source: yamcs/protobuf/events/events_service.proto + +package events + +import ( + api "github.com/nasa/hermes/internal/yamcspb/api" + protoreflect "google.golang.org/protobuf/reflect/protoreflect" + protoimpl "google.golang.org/protobuf/runtime/protoimpl" + timestamppb "google.golang.org/protobuf/types/known/timestamppb" + reflect "reflect" + sync "sync" + unsafe "unsafe" +) + +const ( + // Verify that this generated code is sufficiently up-to-date. + _ = protoimpl.EnforceVersion(20 - protoimpl.MinVersion) + // Verify that runtime/protoimpl is sufficiently up-to-date. + _ = protoimpl.EnforceVersion(protoimpl.MaxVersion - 20) +) + +type ListEventsRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + // Yamcs instance name + Instance *string `protobuf:"bytes,1,opt,name=instance" json:"instance,omitempty"` + // The zero-based row number at which to start outputting results. + // Default: `0` + // + // This option is deprecated and will be removed in a later version. + // Use the returned continuationToken instead. + // + // Deprecated: Marked as deprecated in yamcs/protobuf/events/events_service.proto. + Pos *int64 `protobuf:"varint,2,opt,name=pos" json:"pos,omitempty"` + // The maximum number of returned records per page. Choose this value too high + // and you risk hitting the maximum response size limit enforced by the server. + // Default: `100` + Limit *int32 `protobuf:"varint,3,opt,name=limit" json:"limit,omitempty"` + // The order of the returned results. Can be either `asc` or `desc`. + // Default: `desc` + Order *string `protobuf:"bytes,4,opt,name=order" json:"order,omitempty"` + // The minimum severity level of the events. One of `info`, `watch`, `warning`, + // `distress`, `critical` or `severe`. Default: `info` + Severity *string `protobuf:"bytes,5,opt,name=severity" json:"severity,omitempty"` + // The source of the events. Names must match exactly. + Source []string `protobuf:"bytes,6,rep,name=source" json:"source,omitempty"` + // Continuation token returned by a previous page response. + Next *string `protobuf:"bytes,7,opt,name=next" json:"next,omitempty"` + // Filter the lower bound of the event's generation time. Specify a date string in + // ISO 8601 format. This bound is inclusive. + Start *timestamppb.Timestamp `protobuf:"bytes,8,opt,name=start" json:"start,omitempty"` + // Filter the upper bound of the event's generation time. Specify a date string in + // ISO 8601 format. This bound is exclusive. + Stop *timestamppb.Timestamp `protobuf:"bytes,9,opt,name=stop" json:"stop,omitempty"` + // Text to search for in the message. + Q *string `protobuf:"bytes,10,opt,name=q" json:"q,omitempty"` + // Filter query. See {doc}`../filtering` for how to write a filter query. + // + // Literal text search matches against the `message`, `source` and `type` + // fields. Field comparisons can use any of the following fields: + // + // | Field | Type | Description | + // | ----------- | ------ | ---------------------------------- | + // | `message` | string | Filter by message text | + // | `seqNumber` | number | Filter by archived sequence number | + // | `severity` | enum | Filter by event severity | + // | `source` | string | Filter by event source | + // | `type` | string | Filter by event type | + // + // Possible severities: `info`, `watch`, `warning`, `distress`, `critical` + // or `severe`. + Filter *string `protobuf:"bytes,11,opt,name=filter" json:"filter,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *ListEventsRequest) Reset() { + *x = ListEventsRequest{} + mi := &file_yamcs_protobuf_events_events_service_proto_msgTypes[0] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *ListEventsRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*ListEventsRequest) ProtoMessage() {} + +func (x *ListEventsRequest) ProtoReflect() protoreflect.Message { + mi := &file_yamcs_protobuf_events_events_service_proto_msgTypes[0] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use ListEventsRequest.ProtoReflect.Descriptor instead. +func (*ListEventsRequest) Descriptor() ([]byte, []int) { + return file_yamcs_protobuf_events_events_service_proto_rawDescGZIP(), []int{0} +} + +func (x *ListEventsRequest) GetInstance() string { + if x != nil && x.Instance != nil { + return *x.Instance + } + return "" +} + +// Deprecated: Marked as deprecated in yamcs/protobuf/events/events_service.proto. +func (x *ListEventsRequest) GetPos() int64 { + if x != nil && x.Pos != nil { + return *x.Pos + } + return 0 +} + +func (x *ListEventsRequest) GetLimit() int32 { + if x != nil && x.Limit != nil { + return *x.Limit + } + return 0 +} + +func (x *ListEventsRequest) GetOrder() string { + if x != nil && x.Order != nil { + return *x.Order + } + return "" +} + +func (x *ListEventsRequest) GetSeverity() string { + if x != nil && x.Severity != nil { + return *x.Severity + } + return "" +} + +func (x *ListEventsRequest) GetSource() []string { + if x != nil { + return x.Source + } + return nil +} + +func (x *ListEventsRequest) GetNext() string { + if x != nil && x.Next != nil { + return *x.Next + } + return "" +} + +func (x *ListEventsRequest) GetStart() *timestamppb.Timestamp { + if x != nil { + return x.Start + } + return nil +} + +func (x *ListEventsRequest) GetStop() *timestamppb.Timestamp { + if x != nil { + return x.Stop + } + return nil +} + +func (x *ListEventsRequest) GetQ() string { + if x != nil && x.Q != nil { + return *x.Q + } + return "" +} + +func (x *ListEventsRequest) GetFilter() string { + if x != nil && x.Filter != nil { + return *x.Filter + } + return "" +} + +type ListEventsResponse struct { + state protoimpl.MessageState `protogen:"open.v1"` + // Deprecated, use `events` instead + // + // Deprecated: Marked as deprecated in yamcs/protobuf/events/events_service.proto. + Event []*Event `protobuf:"bytes,1,rep,name=event" json:"event,omitempty"` + // Page with matching events + Events []*Event `protobuf:"bytes,3,rep,name=events" json:"events,omitempty"` + // Token indicating the response is only partial. More results can then + // be obtained by performing the same request (including all original + // query parameters) and setting the `next` parameter to this token. + ContinuationToken *string `protobuf:"bytes,2,opt,name=continuationToken" json:"continuationToken,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *ListEventsResponse) Reset() { + *x = ListEventsResponse{} + mi := &file_yamcs_protobuf_events_events_service_proto_msgTypes[1] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *ListEventsResponse) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*ListEventsResponse) ProtoMessage() {} + +func (x *ListEventsResponse) ProtoReflect() protoreflect.Message { + mi := &file_yamcs_protobuf_events_events_service_proto_msgTypes[1] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use ListEventsResponse.ProtoReflect.Descriptor instead. +func (*ListEventsResponse) Descriptor() ([]byte, []int) { + return file_yamcs_protobuf_events_events_service_proto_rawDescGZIP(), []int{1} +} + +// Deprecated: Marked as deprecated in yamcs/protobuf/events/events_service.proto. +func (x *ListEventsResponse) GetEvent() []*Event { + if x != nil { + return x.Event + } + return nil +} + +func (x *ListEventsResponse) GetEvents() []*Event { + if x != nil { + return x.Events + } + return nil +} + +func (x *ListEventsResponse) GetContinuationToken() string { + if x != nil && x.ContinuationToken != nil { + return *x.ContinuationToken + } + return "" +} + +type SubscribeEventsRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + // Yamcs instance name + Instance *string `protobuf:"bytes,1,opt,name=instance" json:"instance,omitempty"` + // Filter query. See {doc}`../filtering` for how to write a filter query. + // + // Literal text search matches against the `message`, `source` and `type` + // fields. Field comparisons can use any of the following fields: + // + // | Field | Type | Description | + // | ----------- | ------ | ---------------------------------- | + // | `message` | string | Filter by message text | + // | `seqNumber` | number | Filter by archived sequence number | + // | `severity` | enum | Filter by event severity | + // | `source` | string | Filter by event source | + // | `type` | string | Filter by event type | + // + // Possible severities: `info`, `watch`, `warning`, `distress`, `critical` + // or `severe`. + Filter *string `protobuf:"bytes,2,opt,name=filter" json:"filter,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *SubscribeEventsRequest) Reset() { + *x = SubscribeEventsRequest{} + mi := &file_yamcs_protobuf_events_events_service_proto_msgTypes[2] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *SubscribeEventsRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*SubscribeEventsRequest) ProtoMessage() {} + +func (x *SubscribeEventsRequest) ProtoReflect() protoreflect.Message { + mi := &file_yamcs_protobuf_events_events_service_proto_msgTypes[2] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use SubscribeEventsRequest.ProtoReflect.Descriptor instead. +func (*SubscribeEventsRequest) Descriptor() ([]byte, []int) { + return file_yamcs_protobuf_events_events_service_proto_rawDescGZIP(), []int{2} +} + +func (x *SubscribeEventsRequest) GetInstance() string { + if x != nil && x.Instance != nil { + return *x.Instance + } + return "" +} + +func (x *SubscribeEventsRequest) GetFilter() string { + if x != nil && x.Filter != nil { + return *x.Filter + } + return "" +} + +type CreateEventRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + // Yamcs instance name + Instance *string `protobuf:"bytes,1,opt,name=instance" json:"instance,omitempty"` + // Description of the type of the event. Useful for quick classification or filtering. + Type *string `protobuf:"bytes,2,opt,name=type" json:"type,omitempty"` + // **Required.** Event message. + Message *string `protobuf:"bytes,3,opt,name=message" json:"message,omitempty"` + // The severity level of the event. One of `info`, `watch`, `warning`, + // `distress`, `critical` or `severe`. Default is `info` + Severity *string `protobuf:"bytes,4,opt,name=severity" json:"severity,omitempty"` + // Time associated with the event. + // If unspecified, this will default to the current mission time. + Time *timestamppb.Timestamp `protobuf:"bytes,5,opt,name=time" json:"time,omitempty"` + // Source of the event. Useful for grouping events in the archive. Default is + // `User`. + Source *string `protobuf:"bytes,6,opt,name=source" json:"source,omitempty"` + // Sequence number of this event. This is primarily used to determine unicity of + // events coming from the same source. If not set Yamcs will automatically + // assign a sequential number as if every submitted event is unique. + SequenceNumber *int32 `protobuf:"varint,7,opt,name=sequenceNumber" json:"sequenceNumber,omitempty"` + // Additional properties + Extra map[string]string `protobuf:"bytes,8,rep,name=extra" json:"extra,omitempty" protobuf_key:"bytes,1,opt,name=key" protobuf_val:"bytes,2,opt,name=value"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *CreateEventRequest) Reset() { + *x = CreateEventRequest{} + mi := &file_yamcs_protobuf_events_events_service_proto_msgTypes[3] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *CreateEventRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*CreateEventRequest) ProtoMessage() {} + +func (x *CreateEventRequest) ProtoReflect() protoreflect.Message { + mi := &file_yamcs_protobuf_events_events_service_proto_msgTypes[3] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use CreateEventRequest.ProtoReflect.Descriptor instead. +func (*CreateEventRequest) Descriptor() ([]byte, []int) { + return file_yamcs_protobuf_events_events_service_proto_rawDescGZIP(), []int{3} +} + +func (x *CreateEventRequest) GetInstance() string { + if x != nil && x.Instance != nil { + return *x.Instance + } + return "" +} + +func (x *CreateEventRequest) GetType() string { + if x != nil && x.Type != nil { + return *x.Type + } + return "" +} + +func (x *CreateEventRequest) GetMessage() string { + if x != nil && x.Message != nil { + return *x.Message + } + return "" +} + +func (x *CreateEventRequest) GetSeverity() string { + if x != nil && x.Severity != nil { + return *x.Severity + } + return "" +} + +func (x *CreateEventRequest) GetTime() *timestamppb.Timestamp { + if x != nil { + return x.Time + } + return nil +} + +func (x *CreateEventRequest) GetSource() string { + if x != nil && x.Source != nil { + return *x.Source + } + return "" +} + +func (x *CreateEventRequest) GetSequenceNumber() int32 { + if x != nil && x.SequenceNumber != nil { + return *x.SequenceNumber + } + return 0 +} + +func (x *CreateEventRequest) GetExtra() map[string]string { + if x != nil { + return x.Extra + } + return nil +} + +type StreamEventsRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + // Yamcs instance name + Instance *string `protobuf:"bytes,1,opt,name=instance" json:"instance,omitempty"` + // Filter the lower bound of the event's generation time. Specify a date + // string in ISO 8601 format. + Start *timestamppb.Timestamp `protobuf:"bytes,2,opt,name=start" json:"start,omitempty"` + // Filter the upper bound of the event's generation time. Specify a date + // string in ISO 8601 format. + Stop *timestamppb.Timestamp `protobuf:"bytes,3,opt,name=stop" json:"stop,omitempty"` + // Event sources to include. Leave unset, to include all. + Source []string `protobuf:"bytes,4,rep,name=source" json:"source,omitempty"` + // Filter on minimum severity level + Severity *string `protobuf:"bytes,5,opt,name=severity" json:"severity,omitempty"` + // Search by text + Q *string `protobuf:"bytes,6,opt,name=q" json:"q,omitempty"` + // Filter query. See {doc}`../filtering` for how to write a filter query. + // + // Literal text search matches against the `message`, `source` and `type` + // fields. Field comparisons can use any of the following fields: + // + // | Field | Type | Description | + // | ----------- | ------ | ---------------------------------- | + // | `message` | string | Filter by message text | + // | `seqNumber` | number | Filter by archived sequence number | + // | `severity` | enum | Filter by event severity | + // | `source` | string | Filter by event source | + // | `type` | string | Filter by event type | + // + // Possible severities: `info`, `watch`, `warning`, `distress`, `critical` + // or `severe`. + Filter *string `protobuf:"bytes,7,opt,name=filter" json:"filter,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *StreamEventsRequest) Reset() { + *x = StreamEventsRequest{} + mi := &file_yamcs_protobuf_events_events_service_proto_msgTypes[4] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *StreamEventsRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*StreamEventsRequest) ProtoMessage() {} + +func (x *StreamEventsRequest) ProtoReflect() protoreflect.Message { + mi := &file_yamcs_protobuf_events_events_service_proto_msgTypes[4] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use StreamEventsRequest.ProtoReflect.Descriptor instead. +func (*StreamEventsRequest) Descriptor() ([]byte, []int) { + return file_yamcs_protobuf_events_events_service_proto_rawDescGZIP(), []int{4} +} + +func (x *StreamEventsRequest) GetInstance() string { + if x != nil && x.Instance != nil { + return *x.Instance + } + return "" +} + +func (x *StreamEventsRequest) GetStart() *timestamppb.Timestamp { + if x != nil { + return x.Start + } + return nil +} + +func (x *StreamEventsRequest) GetStop() *timestamppb.Timestamp { + if x != nil { + return x.Stop + } + return nil +} + +func (x *StreamEventsRequest) GetSource() []string { + if x != nil { + return x.Source + } + return nil +} + +func (x *StreamEventsRequest) GetSeverity() string { + if x != nil && x.Severity != nil { + return *x.Severity + } + return "" +} + +func (x *StreamEventsRequest) GetQ() string { + if x != nil && x.Q != nil { + return *x.Q + } + return "" +} + +func (x *StreamEventsRequest) GetFilter() string { + if x != nil && x.Filter != nil { + return *x.Filter + } + return "" +} + +type ListEventSourcesRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + // Yamcs instance name + Instance *string `protobuf:"bytes,1,opt,name=instance" json:"instance,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *ListEventSourcesRequest) Reset() { + *x = ListEventSourcesRequest{} + mi := &file_yamcs_protobuf_events_events_service_proto_msgTypes[5] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *ListEventSourcesRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*ListEventSourcesRequest) ProtoMessage() {} + +func (x *ListEventSourcesRequest) ProtoReflect() protoreflect.Message { + mi := &file_yamcs_protobuf_events_events_service_proto_msgTypes[5] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use ListEventSourcesRequest.ProtoReflect.Descriptor instead. +func (*ListEventSourcesRequest) Descriptor() ([]byte, []int) { + return file_yamcs_protobuf_events_events_service_proto_rawDescGZIP(), []int{5} +} + +func (x *ListEventSourcesRequest) GetInstance() string { + if x != nil && x.Instance != nil { + return *x.Instance + } + return "" +} + +type ListEventSourcesResponse struct { + state protoimpl.MessageState `protogen:"open.v1"` + // Deprecated, use `sources` instead + // + // Deprecated: Marked as deprecated in yamcs/protobuf/events/events_service.proto. + Source []string `protobuf:"bytes,1,rep,name=source" json:"source,omitempty"` + // Known event sources + Sources []string `protobuf:"bytes,2,rep,name=sources" json:"sources,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *ListEventSourcesResponse) Reset() { + *x = ListEventSourcesResponse{} + mi := &file_yamcs_protobuf_events_events_service_proto_msgTypes[6] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *ListEventSourcesResponse) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*ListEventSourcesResponse) ProtoMessage() {} + +func (x *ListEventSourcesResponse) ProtoReflect() protoreflect.Message { + mi := &file_yamcs_protobuf_events_events_service_proto_msgTypes[6] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use ListEventSourcesResponse.ProtoReflect.Descriptor instead. +func (*ListEventSourcesResponse) Descriptor() ([]byte, []int) { + return file_yamcs_protobuf_events_events_service_proto_rawDescGZIP(), []int{6} +} + +// Deprecated: Marked as deprecated in yamcs/protobuf/events/events_service.proto. +func (x *ListEventSourcesResponse) GetSource() []string { + if x != nil { + return x.Source + } + return nil +} + +func (x *ListEventSourcesResponse) GetSources() []string { + if x != nil { + return x.Sources + } + return nil +} + +type ExportEventsRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + // Yamcs instance name + Instance *string `protobuf:"bytes,1,opt,name=instance" json:"instance,omitempty"` + // Filter the lower bound of the event's generation time. + // Specify a date string in ISO 8601 format. This bound is inclusive. + Start *timestamppb.Timestamp `protobuf:"bytes,2,opt,name=start" json:"start,omitempty"` + // Filter the upper bound of the event's generation time. Specify a date + // string in ISO 8601 format. This bound is exclusive. + Stop *timestamppb.Timestamp `protobuf:"bytes,3,opt,name=stop" json:"stop,omitempty"` + // The source of the events. Names must match exactly. + Source []string `protobuf:"bytes,4,rep,name=source" json:"source,omitempty"` + // The minimum severity level of the events. One of `info`, `watch`, + // `warning`, `distress` or `severe`. Default: `info` + Severity *string `protobuf:"bytes,5,opt,name=severity" json:"severity,omitempty"` + // Text to search for in the message. + Q *string `protobuf:"bytes,6,opt,name=q" json:"q,omitempty"` + // Filter query. See {doc}`../filtering` for how to write a filter query. + // + // Literal text search matches against the `message`, `source` and `type` + // fields. Field comparisons can use any of the following fields: + // + // | Field | Type | Description | + // | ----------- | ------ | ---------------------------------- | + // | `message` | string | Filter by message text | + // | `seqNumber` | number | Filter by archived sequence number | + // | `severity` | enum | Filter by event severity | + // | `source` | string | Filter by event source | + // | `type` | string | Filter by event type | + // + // Possible severities: `info`, `watch`, `warning`, `distress`, `critical` + // or `severe`. + Filter *string `protobuf:"bytes,8,opt,name=filter" json:"filter,omitempty"` + // Column delimiter. One of `TAB`, `COMMA` or `SEMICOLON`. + // Default: `TAB`. + Delimiter *string `protobuf:"bytes,7,opt,name=delimiter" json:"delimiter,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *ExportEventsRequest) Reset() { + *x = ExportEventsRequest{} + mi := &file_yamcs_protobuf_events_events_service_proto_msgTypes[7] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *ExportEventsRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*ExportEventsRequest) ProtoMessage() {} + +func (x *ExportEventsRequest) ProtoReflect() protoreflect.Message { + mi := &file_yamcs_protobuf_events_events_service_proto_msgTypes[7] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use ExportEventsRequest.ProtoReflect.Descriptor instead. +func (*ExportEventsRequest) Descriptor() ([]byte, []int) { + return file_yamcs_protobuf_events_events_service_proto_rawDescGZIP(), []int{7} +} + +func (x *ExportEventsRequest) GetInstance() string { + if x != nil && x.Instance != nil { + return *x.Instance + } + return "" +} + +func (x *ExportEventsRequest) GetStart() *timestamppb.Timestamp { + if x != nil { + return x.Start + } + return nil +} + +func (x *ExportEventsRequest) GetStop() *timestamppb.Timestamp { + if x != nil { + return x.Stop + } + return nil +} + +func (x *ExportEventsRequest) GetSource() []string { + if x != nil { + return x.Source + } + return nil +} + +func (x *ExportEventsRequest) GetSeverity() string { + if x != nil && x.Severity != nil { + return *x.Severity + } + return "" +} + +func (x *ExportEventsRequest) GetQ() string { + if x != nil && x.Q != nil { + return *x.Q + } + return "" +} + +func (x *ExportEventsRequest) GetFilter() string { + if x != nil && x.Filter != nil { + return *x.Filter + } + return "" +} + +func (x *ExportEventsRequest) GetDelimiter() string { + if x != nil && x.Delimiter != nil { + return *x.Delimiter + } + return "" +} + +var File_yamcs_protobuf_events_events_service_proto protoreflect.FileDescriptor + +const file_yamcs_protobuf_events_events_service_proto_rawDesc = "" + + "\n" + + "*yamcs/protobuf/events/events_service.proto\x12\x15yamcs.protobuf.events\x1a\x1fgoogle/protobuf/timestamp.proto\x1a\x1byamcs/api/annotations.proto\x1a\x18yamcs/api/httpbody.proto\x1a\"yamcs/protobuf/events/events.proto\"\xc1\x02\n" + + "\x11ListEventsRequest\x12\x1a\n" + + "\binstance\x18\x01 \x01(\tR\binstance\x12\x14\n" + + "\x03pos\x18\x02 \x01(\x03B\x02\x18\x01R\x03pos\x12\x14\n" + + "\x05limit\x18\x03 \x01(\x05R\x05limit\x12\x14\n" + + "\x05order\x18\x04 \x01(\tR\x05order\x12\x1a\n" + + "\bseverity\x18\x05 \x01(\tR\bseverity\x12\x16\n" + + "\x06source\x18\x06 \x03(\tR\x06source\x12\x12\n" + + "\x04next\x18\a \x01(\tR\x04next\x120\n" + + "\x05start\x18\b \x01(\v2\x1a.google.protobuf.TimestampR\x05start\x12.\n" + + "\x04stop\x18\t \x01(\v2\x1a.google.protobuf.TimestampR\x04stop\x12\f\n" + + "\x01q\x18\n" + + " \x01(\tR\x01q\x12\x16\n" + + "\x06filter\x18\v \x01(\tR\x06filter\"\xb0\x01\n" + + "\x12ListEventsResponse\x126\n" + + "\x05event\x18\x01 \x03(\v2\x1c.yamcs.protobuf.events.EventB\x02\x18\x01R\x05event\x124\n" + + "\x06events\x18\x03 \x03(\v2\x1c.yamcs.protobuf.events.EventR\x06events\x12,\n" + + "\x11continuationToken\x18\x02 \x01(\tR\x11continuationToken\"L\n" + + "\x16SubscribeEventsRequest\x12\x1a\n" + + "\binstance\x18\x01 \x01(\tR\binstance\x12\x16\n" + + "\x06filter\x18\x02 \x01(\tR\x06filter\"\xf0\x02\n" + + "\x12CreateEventRequest\x12\x1a\n" + + "\binstance\x18\x01 \x01(\tR\binstance\x12\x12\n" + + "\x04type\x18\x02 \x01(\tR\x04type\x12\x18\n" + + "\amessage\x18\x03 \x01(\tR\amessage\x12\x1a\n" + + "\bseverity\x18\x04 \x01(\tR\bseverity\x12.\n" + + "\x04time\x18\x05 \x01(\v2\x1a.google.protobuf.TimestampR\x04time\x12\x16\n" + + "\x06source\x18\x06 \x01(\tR\x06source\x12&\n" + + "\x0esequenceNumber\x18\a \x01(\x05R\x0esequenceNumber\x12J\n" + + "\x05extra\x18\b \x03(\v24.yamcs.protobuf.events.CreateEventRequest.ExtraEntryR\x05extra\x1a8\n" + + "\n" + + "ExtraEntry\x12\x10\n" + + "\x03key\x18\x01 \x01(\tR\x03key\x12\x14\n" + + "\x05value\x18\x02 \x01(\tR\x05value:\x028\x01\"\xed\x01\n" + + "\x13StreamEventsRequest\x12\x1a\n" + + "\binstance\x18\x01 \x01(\tR\binstance\x120\n" + + "\x05start\x18\x02 \x01(\v2\x1a.google.protobuf.TimestampR\x05start\x12.\n" + + "\x04stop\x18\x03 \x01(\v2\x1a.google.protobuf.TimestampR\x04stop\x12\x16\n" + + "\x06source\x18\x04 \x03(\tR\x06source\x12\x1a\n" + + "\bseverity\x18\x05 \x01(\tR\bseverity\x12\f\n" + + "\x01q\x18\x06 \x01(\tR\x01q\x12\x16\n" + + "\x06filter\x18\a \x01(\tR\x06filter\"5\n" + + "\x17ListEventSourcesRequest\x12\x1a\n" + + "\binstance\x18\x01 \x01(\tR\binstance\"P\n" + + "\x18ListEventSourcesResponse\x12\x1a\n" + + "\x06source\x18\x01 \x03(\tB\x02\x18\x01R\x06source\x12\x18\n" + + "\asources\x18\x02 \x03(\tR\asources\"\x8b\x02\n" + + "\x13ExportEventsRequest\x12\x1a\n" + + "\binstance\x18\x01 \x01(\tR\binstance\x120\n" + + "\x05start\x18\x02 \x01(\v2\x1a.google.protobuf.TimestampR\x05start\x12.\n" + + "\x04stop\x18\x03 \x01(\v2\x1a.google.protobuf.TimestampR\x04stop\x12\x16\n" + + "\x06source\x18\x04 \x03(\tR\x06source\x12\x1a\n" + + "\bseverity\x18\x05 \x01(\tR\bseverity\x12\f\n" + + "\x01q\x18\x06 \x01(\tR\x01q\x12\x16\n" + + "\x06filter\x18\b \x01(\tR\x06filter\x12\x1c\n" + + "\tdelimiter\x18\a \x01(\tR\tdelimiter2\xee\x06\n" + + "\tEventsApi\x12\xb1\x01\n" + + "\n" + + "ListEvents\x12(.yamcs.protobuf.events.ListEventsRequest\x1a).yamcs.protobuf.events.ListEventsResponse\"N\x8a\x92\x03JZ(:\x01*\x1a#/api/archive/{instance}/events:list\n" + + "\x1e/api/archive/{instance}/events\x12\x7f\n" + + "\vCreateEvent\x12).yamcs.protobuf.events.CreateEventRequest\x1a\x1c.yamcs.protobuf.events.Event\"'\x8a\x92\x03#:\x01*\x1a\x1e/api/archive/{instance}/events\x12\xa1\x01\n" + + "\x10ListEventSources\x12..yamcs.protobuf.events.ListEventSourcesRequest\x1a/.yamcs.protobuf.events.ListEventSourcesResponse\",\x8a\x92\x03(\n" + + "&/api/archive/{instance}/events/sources\x12\x90\x01\n" + + "\fStreamEvents\x12*.yamcs.protobuf.events.StreamEventsRequest\x1a\x1c.yamcs.protobuf.events.Event\"4\x8a\x92\x030:\x01*\x1a+/api/stream-archive/{instance}:streamEvents0\x01\x12}\n" + + "\fExportEvents\x12*.yamcs.protobuf.events.ExportEventsRequest\x1a\x13.yamcs.api.HttpBody\"*\x8a\x92\x03&\n" + + "$/api/archive/{instance}:exportEvents0\x01\x12p\n" + + "\x0fSubscribeEvents\x12-.yamcs.protobuf.events.SubscribeEventsRequest\x1a\x1c.yamcs.protobuf.events.Event\"\fڒ\x03\b\n" + + "\x06events(\x010\x01\x1a\x04Ѐ\x01\x01B*\n" + + "\x12org.yamcs.protobufB\x12EventsServiceProtoP\x01" + +var ( + file_yamcs_protobuf_events_events_service_proto_rawDescOnce sync.Once + file_yamcs_protobuf_events_events_service_proto_rawDescData []byte +) + +func file_yamcs_protobuf_events_events_service_proto_rawDescGZIP() []byte { + file_yamcs_protobuf_events_events_service_proto_rawDescOnce.Do(func() { + file_yamcs_protobuf_events_events_service_proto_rawDescData = protoimpl.X.CompressGZIP(unsafe.Slice(unsafe.StringData(file_yamcs_protobuf_events_events_service_proto_rawDesc), len(file_yamcs_protobuf_events_events_service_proto_rawDesc))) + }) + return file_yamcs_protobuf_events_events_service_proto_rawDescData +} + +var file_yamcs_protobuf_events_events_service_proto_msgTypes = make([]protoimpl.MessageInfo, 9) +var file_yamcs_protobuf_events_events_service_proto_goTypes = []any{ + (*ListEventsRequest)(nil), // 0: yamcs.protobuf.events.ListEventsRequest + (*ListEventsResponse)(nil), // 1: yamcs.protobuf.events.ListEventsResponse + (*SubscribeEventsRequest)(nil), // 2: yamcs.protobuf.events.SubscribeEventsRequest + (*CreateEventRequest)(nil), // 3: yamcs.protobuf.events.CreateEventRequest + (*StreamEventsRequest)(nil), // 4: yamcs.protobuf.events.StreamEventsRequest + (*ListEventSourcesRequest)(nil), // 5: yamcs.protobuf.events.ListEventSourcesRequest + (*ListEventSourcesResponse)(nil), // 6: yamcs.protobuf.events.ListEventSourcesResponse + (*ExportEventsRequest)(nil), // 7: yamcs.protobuf.events.ExportEventsRequest + nil, // 8: yamcs.protobuf.events.CreateEventRequest.ExtraEntry + (*timestamppb.Timestamp)(nil), // 9: google.protobuf.Timestamp + (*Event)(nil), // 10: yamcs.protobuf.events.Event + (*api.HttpBody)(nil), // 11: yamcs.api.HttpBody +} +var file_yamcs_protobuf_events_events_service_proto_depIdxs = []int32{ + 9, // 0: yamcs.protobuf.events.ListEventsRequest.start:type_name -> google.protobuf.Timestamp + 9, // 1: yamcs.protobuf.events.ListEventsRequest.stop:type_name -> google.protobuf.Timestamp + 10, // 2: yamcs.protobuf.events.ListEventsResponse.event:type_name -> yamcs.protobuf.events.Event + 10, // 3: yamcs.protobuf.events.ListEventsResponse.events:type_name -> yamcs.protobuf.events.Event + 9, // 4: yamcs.protobuf.events.CreateEventRequest.time:type_name -> google.protobuf.Timestamp + 8, // 5: yamcs.protobuf.events.CreateEventRequest.extra:type_name -> yamcs.protobuf.events.CreateEventRequest.ExtraEntry + 9, // 6: yamcs.protobuf.events.StreamEventsRequest.start:type_name -> google.protobuf.Timestamp + 9, // 7: yamcs.protobuf.events.StreamEventsRequest.stop:type_name -> google.protobuf.Timestamp + 9, // 8: yamcs.protobuf.events.ExportEventsRequest.start:type_name -> google.protobuf.Timestamp + 9, // 9: yamcs.protobuf.events.ExportEventsRequest.stop:type_name -> google.protobuf.Timestamp + 0, // 10: yamcs.protobuf.events.EventsApi.ListEvents:input_type -> yamcs.protobuf.events.ListEventsRequest + 3, // 11: yamcs.protobuf.events.EventsApi.CreateEvent:input_type -> yamcs.protobuf.events.CreateEventRequest + 5, // 12: yamcs.protobuf.events.EventsApi.ListEventSources:input_type -> yamcs.protobuf.events.ListEventSourcesRequest + 4, // 13: yamcs.protobuf.events.EventsApi.StreamEvents:input_type -> yamcs.protobuf.events.StreamEventsRequest + 7, // 14: yamcs.protobuf.events.EventsApi.ExportEvents:input_type -> yamcs.protobuf.events.ExportEventsRequest + 2, // 15: yamcs.protobuf.events.EventsApi.SubscribeEvents:input_type -> yamcs.protobuf.events.SubscribeEventsRequest + 1, // 16: yamcs.protobuf.events.EventsApi.ListEvents:output_type -> yamcs.protobuf.events.ListEventsResponse + 10, // 17: yamcs.protobuf.events.EventsApi.CreateEvent:output_type -> yamcs.protobuf.events.Event + 6, // 18: yamcs.protobuf.events.EventsApi.ListEventSources:output_type -> yamcs.protobuf.events.ListEventSourcesResponse + 10, // 19: yamcs.protobuf.events.EventsApi.StreamEvents:output_type -> yamcs.protobuf.events.Event + 11, // 20: yamcs.protobuf.events.EventsApi.ExportEvents:output_type -> yamcs.api.HttpBody + 10, // 21: yamcs.protobuf.events.EventsApi.SubscribeEvents:output_type -> yamcs.protobuf.events.Event + 16, // [16:22] is the sub-list for method output_type + 10, // [10:16] is the sub-list for method input_type + 10, // [10:10] is the sub-list for extension type_name + 10, // [10:10] is the sub-list for extension extendee + 0, // [0:10] is the sub-list for field type_name +} + +func init() { file_yamcs_protobuf_events_events_service_proto_init() } +func file_yamcs_protobuf_events_events_service_proto_init() { + if File_yamcs_protobuf_events_events_service_proto != nil { + return + } + file_yamcs_protobuf_events_events_proto_init() + type x struct{} + out := protoimpl.TypeBuilder{ + File: protoimpl.DescBuilder{ + GoPackagePath: reflect.TypeOf(x{}).PkgPath(), + RawDescriptor: unsafe.Slice(unsafe.StringData(file_yamcs_protobuf_events_events_service_proto_rawDesc), len(file_yamcs_protobuf_events_events_service_proto_rawDesc)), + NumEnums: 0, + NumMessages: 9, + NumExtensions: 0, + NumServices: 1, + }, + GoTypes: file_yamcs_protobuf_events_events_service_proto_goTypes, + DependencyIndexes: file_yamcs_protobuf_events_events_service_proto_depIdxs, + MessageInfos: file_yamcs_protobuf_events_events_service_proto_msgTypes, + }.Build() + File_yamcs_protobuf_events_events_service_proto = out.File + file_yamcs_protobuf_events_events_service_proto_goTypes = nil + file_yamcs_protobuf_events_events_service_proto_depIdxs = nil +} diff --git a/internal/yamcspb/protobuf/events/events_service_grpc.pb.go b/internal/yamcspb/protobuf/events/events_service_grpc.pb.go new file mode 100644 index 00000000..69eef7be --- /dev/null +++ b/internal/yamcspb/protobuf/events/events_service_grpc.pb.go @@ -0,0 +1,331 @@ +// Code generated by protoc-gen-go-grpc. DO NOT EDIT. +// versions: +// - protoc-gen-go-grpc v1.6.2 +// - protoc v7.36.0 +// source: yamcs/protobuf/events/events_service.proto + +package events + +import ( + context "context" + api "github.com/nasa/hermes/internal/yamcspb/api" + grpc "google.golang.org/grpc" + codes "google.golang.org/grpc/codes" + status "google.golang.org/grpc/status" +) + +// This is a compile-time assertion to ensure that this generated file +// is compatible with the grpc package it is being compiled against. +// Requires gRPC-Go v1.64.0 or later. +const _ = grpc.SupportPackageIsVersion9 + +const ( + EventsApi_ListEvents_FullMethodName = "/yamcs.protobuf.events.EventsApi/ListEvents" + EventsApi_CreateEvent_FullMethodName = "/yamcs.protobuf.events.EventsApi/CreateEvent" + EventsApi_ListEventSources_FullMethodName = "/yamcs.protobuf.events.EventsApi/ListEventSources" + EventsApi_StreamEvents_FullMethodName = "/yamcs.protobuf.events.EventsApi/StreamEvents" + EventsApi_ExportEvents_FullMethodName = "/yamcs.protobuf.events.EventsApi/ExportEvents" + EventsApi_SubscribeEvents_FullMethodName = "/yamcs.protobuf.events.EventsApi/SubscribeEvents" +) + +// EventsApiClient is the client API for EventsApi service. +// +// For semantics around ctx use and closing/ending streaming RPCs, please refer to https://pkg.go.dev/google.golang.org/grpc/?tab=doc#ClientConn.NewStream. +type EventsApiClient interface { + // List events + // + // Two alternative URI forms can be used. The `GET` method is convenient, while + // the `POST` method allows passing long filter expressions. + ListEvents(ctx context.Context, in *ListEventsRequest, opts ...grpc.CallOption) (*ListEventsResponse, error) + // Create an event + CreateEvent(ctx context.Context, in *CreateEventRequest, opts ...grpc.CallOption) (*Event, error) + // List event sources + ListEventSources(ctx context.Context, in *ListEventSourcesRequest, opts ...grpc.CallOption) (*ListEventSourcesResponse, error) + // Streams back events + StreamEvents(ctx context.Context, in *StreamEventsRequest, opts ...grpc.CallOption) (grpc.ServerStreamingClient[Event], error) + // Export events in CSV format + ExportEvents(ctx context.Context, in *ExportEventsRequest, opts ...grpc.CallOption) (grpc.ServerStreamingClient[api.HttpBody], error) + // Receive event updates + SubscribeEvents(ctx context.Context, opts ...grpc.CallOption) (grpc.BidiStreamingClient[SubscribeEventsRequest, Event], error) +} + +type eventsApiClient struct { + cc grpc.ClientConnInterface +} + +func NewEventsApiClient(cc grpc.ClientConnInterface) EventsApiClient { + return &eventsApiClient{cc} +} + +func (c *eventsApiClient) ListEvents(ctx context.Context, in *ListEventsRequest, opts ...grpc.CallOption) (*ListEventsResponse, error) { + cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) + out := new(ListEventsResponse) + err := c.cc.Invoke(ctx, EventsApi_ListEvents_FullMethodName, in, out, cOpts...) + if err != nil { + return nil, err + } + return out, nil +} + +func (c *eventsApiClient) CreateEvent(ctx context.Context, in *CreateEventRequest, opts ...grpc.CallOption) (*Event, error) { + cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) + out := new(Event) + err := c.cc.Invoke(ctx, EventsApi_CreateEvent_FullMethodName, in, out, cOpts...) + if err != nil { + return nil, err + } + return out, nil +} + +func (c *eventsApiClient) ListEventSources(ctx context.Context, in *ListEventSourcesRequest, opts ...grpc.CallOption) (*ListEventSourcesResponse, error) { + cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) + out := new(ListEventSourcesResponse) + err := c.cc.Invoke(ctx, EventsApi_ListEventSources_FullMethodName, in, out, cOpts...) + if err != nil { + return nil, err + } + return out, nil +} + +func (c *eventsApiClient) StreamEvents(ctx context.Context, in *StreamEventsRequest, opts ...grpc.CallOption) (grpc.ServerStreamingClient[Event], error) { + cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) + stream, err := c.cc.NewStream(ctx, &EventsApi_ServiceDesc.Streams[0], EventsApi_StreamEvents_FullMethodName, cOpts...) + if err != nil { + return nil, err + } + x := &grpc.GenericClientStream[StreamEventsRequest, Event]{ClientStream: stream} + if err := x.ClientStream.SendMsg(in); err != nil { + return nil, err + } + if err := x.ClientStream.CloseSend(); err != nil { + return nil, err + } + return x, nil +} + +// This type alias is provided for backwards compatibility with existing code that references the prior non-generic stream type by name. +type EventsApi_StreamEventsClient = grpc.ServerStreamingClient[Event] + +func (c *eventsApiClient) ExportEvents(ctx context.Context, in *ExportEventsRequest, opts ...grpc.CallOption) (grpc.ServerStreamingClient[api.HttpBody], error) { + cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) + stream, err := c.cc.NewStream(ctx, &EventsApi_ServiceDesc.Streams[1], EventsApi_ExportEvents_FullMethodName, cOpts...) + if err != nil { + return nil, err + } + x := &grpc.GenericClientStream[ExportEventsRequest, api.HttpBody]{ClientStream: stream} + if err := x.ClientStream.SendMsg(in); err != nil { + return nil, err + } + if err := x.ClientStream.CloseSend(); err != nil { + return nil, err + } + return x, nil +} + +// This type alias is provided for backwards compatibility with existing code that references the prior non-generic stream type by name. +type EventsApi_ExportEventsClient = grpc.ServerStreamingClient[api.HttpBody] + +func (c *eventsApiClient) SubscribeEvents(ctx context.Context, opts ...grpc.CallOption) (grpc.BidiStreamingClient[SubscribeEventsRequest, Event], error) { + cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) + stream, err := c.cc.NewStream(ctx, &EventsApi_ServiceDesc.Streams[2], EventsApi_SubscribeEvents_FullMethodName, cOpts...) + if err != nil { + return nil, err + } + x := &grpc.GenericClientStream[SubscribeEventsRequest, Event]{ClientStream: stream} + return x, nil +} + +// This type alias is provided for backwards compatibility with existing code that references the prior non-generic stream type by name. +type EventsApi_SubscribeEventsClient = grpc.BidiStreamingClient[SubscribeEventsRequest, Event] + +// EventsApiServer is the server API for EventsApi service. +// All implementations must embed UnimplementedEventsApiServer +// for forward compatibility. +type EventsApiServer interface { + // List events + // + // Two alternative URI forms can be used. The `GET` method is convenient, while + // the `POST` method allows passing long filter expressions. + ListEvents(context.Context, *ListEventsRequest) (*ListEventsResponse, error) + // Create an event + CreateEvent(context.Context, *CreateEventRequest) (*Event, error) + // List event sources + ListEventSources(context.Context, *ListEventSourcesRequest) (*ListEventSourcesResponse, error) + // Streams back events + StreamEvents(*StreamEventsRequest, grpc.ServerStreamingServer[Event]) error + // Export events in CSV format + ExportEvents(*ExportEventsRequest, grpc.ServerStreamingServer[api.HttpBody]) error + // Receive event updates + SubscribeEvents(grpc.BidiStreamingServer[SubscribeEventsRequest, Event]) error + mustEmbedUnimplementedEventsApiServer() +} + +// UnimplementedEventsApiServer must be embedded to have +// forward compatible implementations. +// +// NOTE: this should be embedded by value instead of pointer to avoid a nil +// pointer dereference when methods are called. +type UnimplementedEventsApiServer struct{} + +func (UnimplementedEventsApiServer) ListEvents(context.Context, *ListEventsRequest) (*ListEventsResponse, error) { + return nil, status.Error(codes.Unimplemented, "method ListEvents not implemented") +} +func (UnimplementedEventsApiServer) CreateEvent(context.Context, *CreateEventRequest) (*Event, error) { + return nil, status.Error(codes.Unimplemented, "method CreateEvent not implemented") +} +func (UnimplementedEventsApiServer) ListEventSources(context.Context, *ListEventSourcesRequest) (*ListEventSourcesResponse, error) { + return nil, status.Error(codes.Unimplemented, "method ListEventSources not implemented") +} +func (UnimplementedEventsApiServer) StreamEvents(*StreamEventsRequest, grpc.ServerStreamingServer[Event]) error { + return status.Error(codes.Unimplemented, "method StreamEvents not implemented") +} +func (UnimplementedEventsApiServer) ExportEvents(*ExportEventsRequest, grpc.ServerStreamingServer[api.HttpBody]) error { + return status.Error(codes.Unimplemented, "method ExportEvents not implemented") +} +func (UnimplementedEventsApiServer) SubscribeEvents(grpc.BidiStreamingServer[SubscribeEventsRequest, Event]) error { + return status.Error(codes.Unimplemented, "method SubscribeEvents not implemented") +} +func (UnimplementedEventsApiServer) mustEmbedUnimplementedEventsApiServer() {} +func (UnimplementedEventsApiServer) testEmbeddedByValue() {} + +// UnsafeEventsApiServer may be embedded to opt out of forward compatibility for this service. +// Use of this interface is not recommended, as added methods to EventsApiServer will +// result in compilation errors. +type UnsafeEventsApiServer interface { + mustEmbedUnimplementedEventsApiServer() +} + +func RegisterEventsApiServer(s grpc.ServiceRegistrar, srv EventsApiServer) { + // If the following call panics, it indicates UnimplementedEventsApiServer was + // embedded by pointer and is nil. This will cause panics if an + // unimplemented method is ever invoked, so we test this at initialization + // time to prevent it from happening at runtime later due to I/O. + if t, ok := srv.(interface{ testEmbeddedByValue() }); ok { + t.testEmbeddedByValue() + } + s.RegisterService(&EventsApi_ServiceDesc, srv) +} + +func _EventsApi_ListEvents_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { + in := new(ListEventsRequest) + if err := dec(in); err != nil { + return nil, err + } + if interceptor == nil { + return srv.(EventsApiServer).ListEvents(ctx, in) + } + info := &grpc.UnaryServerInfo{ + Server: srv, + FullMethod: EventsApi_ListEvents_FullMethodName, + } + handler := func(ctx context.Context, req interface{}) (interface{}, error) { + return srv.(EventsApiServer).ListEvents(ctx, req.(*ListEventsRequest)) + } + return interceptor(ctx, in, info, handler) +} + +func _EventsApi_CreateEvent_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { + in := new(CreateEventRequest) + if err := dec(in); err != nil { + return nil, err + } + if interceptor == nil { + return srv.(EventsApiServer).CreateEvent(ctx, in) + } + info := &grpc.UnaryServerInfo{ + Server: srv, + FullMethod: EventsApi_CreateEvent_FullMethodName, + } + handler := func(ctx context.Context, req interface{}) (interface{}, error) { + return srv.(EventsApiServer).CreateEvent(ctx, req.(*CreateEventRequest)) + } + return interceptor(ctx, in, info, handler) +} + +func _EventsApi_ListEventSources_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { + in := new(ListEventSourcesRequest) + if err := dec(in); err != nil { + return nil, err + } + if interceptor == nil { + return srv.(EventsApiServer).ListEventSources(ctx, in) + } + info := &grpc.UnaryServerInfo{ + Server: srv, + FullMethod: EventsApi_ListEventSources_FullMethodName, + } + handler := func(ctx context.Context, req interface{}) (interface{}, error) { + return srv.(EventsApiServer).ListEventSources(ctx, req.(*ListEventSourcesRequest)) + } + return interceptor(ctx, in, info, handler) +} + +func _EventsApi_StreamEvents_Handler(srv interface{}, stream grpc.ServerStream) error { + m := new(StreamEventsRequest) + if err := stream.RecvMsg(m); err != nil { + return err + } + return srv.(EventsApiServer).StreamEvents(m, &grpc.GenericServerStream[StreamEventsRequest, Event]{ServerStream: stream}) +} + +// This type alias is provided for backwards compatibility with existing code that references the prior non-generic stream type by name. +type EventsApi_StreamEventsServer = grpc.ServerStreamingServer[Event] + +func _EventsApi_ExportEvents_Handler(srv interface{}, stream grpc.ServerStream) error { + m := new(ExportEventsRequest) + if err := stream.RecvMsg(m); err != nil { + return err + } + return srv.(EventsApiServer).ExportEvents(m, &grpc.GenericServerStream[ExportEventsRequest, api.HttpBody]{ServerStream: stream}) +} + +// This type alias is provided for backwards compatibility with existing code that references the prior non-generic stream type by name. +type EventsApi_ExportEventsServer = grpc.ServerStreamingServer[api.HttpBody] + +func _EventsApi_SubscribeEvents_Handler(srv interface{}, stream grpc.ServerStream) error { + return srv.(EventsApiServer).SubscribeEvents(&grpc.GenericServerStream[SubscribeEventsRequest, Event]{ServerStream: stream}) +} + +// This type alias is provided for backwards compatibility with existing code that references the prior non-generic stream type by name. +type EventsApi_SubscribeEventsServer = grpc.BidiStreamingServer[SubscribeEventsRequest, Event] + +// EventsApi_ServiceDesc is the grpc.ServiceDesc for EventsApi service. +// It's only intended for direct use with grpc.RegisterService, +// and not to be introspected or modified (even as a copy) +var EventsApi_ServiceDesc = grpc.ServiceDesc{ + ServiceName: "yamcs.protobuf.events.EventsApi", + HandlerType: (*EventsApiServer)(nil), + Methods: []grpc.MethodDesc{ + { + MethodName: "ListEvents", + Handler: _EventsApi_ListEvents_Handler, + }, + { + MethodName: "CreateEvent", + Handler: _EventsApi_CreateEvent_Handler, + }, + { + MethodName: "ListEventSources", + Handler: _EventsApi_ListEventSources_Handler, + }, + }, + Streams: []grpc.StreamDesc{ + { + StreamName: "StreamEvents", + Handler: _EventsApi_StreamEvents_Handler, + ServerStreams: true, + }, + { + StreamName: "ExportEvents", + Handler: _EventsApi_ExportEvents_Handler, + ServerStreams: true, + }, + { + StreamName: "SubscribeEvents", + Handler: _EventsApi_SubscribeEvents_Handler, + ServerStreams: true, + ClientStreams: true, + }, + }, + Metadata: "yamcs/protobuf/events/events_service.proto", +}