master
go 73 lines 1.73 KB
Raw
1 // SPDX-License-Identifier: GPL-3.0-or-later
2
3 package main
4
5 import (
6 "context"
7 "flag"
8 "fmt"
9 "os"
10 "strings"
11 "time"
12
13 "gopkg.in/rethinkdb/rethinkdb-go.v6"
14 )
15
16 func main() {
17 addr := flag.String("addr", "127.0.0.1:28015", "RethinkDB address")
18 duration := flag.Duration("duration", 30*time.Second, "How long to keep the changefeed open")
19 flag.Parse()
20
21 sess, err := rethinkdb.Connect(rethinkdb.ConnectOpts{Address: *addr})
22 if err != nil {
23 fatalf("connect: %v", err)
24 }
25 defer func() { _ = sess.Close() }()
26
27 ensureDB(sess, "netdata")
28 ensureTable(sess, "netdata", "demo")
29 insertSeed(sess, "netdata", "demo")
30
31 ctx, cancel := context.WithTimeout(context.Background(), *duration)
32 defer cancel()
33
34 cur, err := rethinkdb.DB("netdata").Table("demo").Changes().Run(sess, rethinkdb.RunOpts{Context: ctx})
35 if err != nil {
36 fatalf("start changefeed: %v", err)
37 }
38 defer func() { _ = cur.Close() }()
39
40 <-ctx.Done()
41 }
42
43 func ensureDB(sess *rethinkdb.Session, name string) {
44 if _, err := rethinkdb.DBCreate(name).RunWrite(sess); err != nil {
45 if !isAlreadyExists(err) {
46 fatalf("db create: %v", err)
47 }
48 }
49 }
50
51 func ensureTable(sess *rethinkdb.Session, db, table string) {
52 if _, err := rethinkdb.DB(db).TableCreate(table).RunWrite(sess); err != nil {
53 if !isAlreadyExists(err) {
54 fatalf("table create: %v", err)
55 }
56 }
57 }
58
59 func insertSeed(sess *rethinkdb.Session, db, table string) {
60 _, _ = rethinkdb.DB(db).Table(table).Insert(map[string]any{
61 "id": "seed",
62 "name": "alpha",
63 }).RunWrite(sess)
64 }
65
66 func isAlreadyExists(err error) bool {
67 return err != nil && (strings.Contains(err.Error(), "already exists") || strings.Contains(err.Error(), "Duplicate"))
68 }
69
70 func fatalf(format string, args ...any) {
71 _, _ = fmt.Fprintf(os.Stderr, format+"\n", args...)
72 os.Exit(1)
73 }