Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
61 changes: 55 additions & 6 deletions cmd/yamcs-recorder/README.md
Original file line number Diff line number Diff line change
@@ -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
Expand All @@ -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.<key>`; 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' \
Expand Down
150 changes: 139 additions & 11 deletions cmd/yamcs-recorder/main.go
Original file line number Diff line number Diff line change
@@ -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 (
Expand All @@ -10,22 +12,32 @@ import (
"errors"
"fmt"
"log/slog"
"net"
"net/url"
"os"
"os/signal"
"slices"
"strconv"
"strings"
"sync/atomic"
"syscall"
"time"

"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"
)

Expand All @@ -36,6 +48,7 @@ type config struct {
instance string
processor string
postgresql string
otlp string
includeTopLevel bool
}

Expand All @@ -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")
Expand All @@ -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.
Expand All @@ -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
Expand All @@ -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,
Expand All @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand All @@ -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.
Expand Down
3 changes: 2 additions & 1 deletion cmd/yamcs-recorder/main_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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")
}
22 changes: 22 additions & 0 deletions docker-compose.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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:
2 changes: 1 addition & 1 deletion grafana-datasource-plugin/src/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.

<br>

Expand Down
2 changes: 1 addition & 1 deletion internal/recorder/convert/convert.go
Original file line number Diff line number Diff line change
@@ -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 (
Expand Down
Loading
Loading