package link import ( "context" "encoding/json" "errors" "fmt" "sort" "github.com/nats-io/nats.go" "github.com/nats-io/nats.go/jetstream" "github.com/novox/mesh-controller/internal/broker" ) // Calls kept on the bus (novox/hq to-be 45 ยง6). // // **Two keys a call**: `` holds the record โ€” verb, arguments as kept, caller, state, times โ€” and // `.answer` the answer, bounded. A listing watches the records alone, so reading the last thousand // calls reads the last thousand small records and not a thousand answers; one call asked by id reads // both. A call's id is one token, so the record's key never holds a dot and `*` matches records only. // BusCalls keeps calls in the controller's calls bucket. type BusCalls struct { kv jetstream.KeyValue } // CallsOnTheBus opens the calls bucket the controller asserts at its start. func CallsOnTheBus(ctx context.Context, conn *nats.Conn) (*BusCalls, error) { api, err := jetstream.New(conn) if err != nil { return nil, err } kv, err := api.KeyValue(ctx, broker.CallsBucket) if err != nil { return nil, fmt.Errorf("the calls bucket %s is not on the bus โ€” the controller asserts it at its "+ "start, so one older than this has not: %w", broker.CallsBucket, err) } return &BusCalls{kv: kv}, nil } const answerKey = ".answer" // Keep writes a call's record, and its answer when it has one. func (b *BusCalls) Keep(ctx context.Context, c Call) error { answer := c.Answer c.Answer = nil if len(answer) > 0 { if len(answer) > broker.CallAnswerBytes { answer, _ = json.Marshal(map[string]any{"cut": fmt.Sprintf("an answer of %d bytes; the first %d "+ "are kept", len(answer), broker.CallAnswerBytes), "start": string(answer[:broker.CallAnswerBytes])}) } // The answer before the record, so a record saying a call finished never points at an answer // not yet written. if _, err := b.kv.Put(ctx, c.ID+answerKey, answer); err != nil { return err } } record, err := json.Marshal(c) if err != nil { return err } _, err = b.kv.Put(ctx, c.ID, record) return err } // Kept is one call, with its answer. func (b *BusCalls) Kept(ctx context.Context, id string) (Call, bool, error) { entry, err := b.kv.Get(ctx, id) if errors.Is(err, jetstream.ErrKeyNotFound) || errors.Is(err, jetstream.ErrInvalidKey) { return Call{}, false, nil } if err != nil { return Call{}, false, err } var c Call if err := json.Unmarshal(entry.Value(), &c); err != nil { return Call{}, false, fmt.Errorf("the record of %s on the bus is not a call: %w", id, err) } answer, err := b.kv.Get(ctx, id+answerKey) switch { case err == nil: c.Answer = json.RawMessage(answer.Value()) case !errors.Is(err, jetstream.ErrKeyNotFound): return Call{}, false, err } return c, true, nil } // Recent is every call the bucket holds, newest first, without answers: read once through a watch // of the records, which hands over the current value of each and then says it has. func (b *BusCalls) Recent(ctx context.Context) ([]Call, error) { w, err := b.kv.Watch(ctx, "*", jetstream.IgnoreDeletes()) if err != nil { return nil, err } defer func() { _ = w.Stop() }() var out []Call for { select { case <-ctx.Done(): return nil, fmt.Errorf("reading the calls bucket: %w", ctx.Err()) case entry := <-w.Updates(): if entry == nil { // Every current value handed over. sort.SliceStable(out, func(i, j int) bool { return out[i].Started.After(out[j].Started) }) return out, nil } var c Call if json.Unmarshal(entry.Value(), &c) == nil && c.ID != "" { out = append(out, c) } } } }