diff --git a/knotserver/config/config.go b/knotserver/config/config.go index 116bb6ea..bb218173 100644 --- a/knotserver/config/config.go +++ b/knotserver/config/config.go @@ -22,8 +22,9 @@ type Server struct { } type Config struct { - Repo Repo `env:",prefix=KNOT_REPO_"` - Server Server `env:",prefix=KNOT_SERVER_"` + Repo Repo `env:",prefix=KNOT_REPO_"` + Server Server `env:",prefix=KNOT_SERVER_"` + AppViewEndpoint string `env:"APPVIEW_ENDPOINT, default=https://tangled.sh"` } func Load(ctx context.Context) (*Config, error) { diff --git a/knotserver/handler.go b/knotserver/handler.go index 5feacef9..bc02182f 100644 --- a/knotserver/handler.go +++ b/knotserver/handler.go @@ -2,14 +2,10 @@ package knotserver import ( "context" - "encoding/json" "fmt" - "log" "net/http" - "time" "github.com/go-chi/chi/v5" - tangled "github.com/sotangled/tangled/api/tangled" "github.com/sotangled/tangled/knotserver/config" "github.com/sotangled/tangled/knotserver/db" "github.com/sotangled/tangled/knotserver/jsclient" @@ -94,6 +90,11 @@ func Setup(ctx context.Context, c *config.Config, db *db.DB, e *rbac.Enforcer) ( r.Put("/new", h.NewRepo) }) + r.Route("/member", func(r chi.Router) { + r.Use(h.VerifySignature) + r.Put("/add", h.NewRepo) + }) + // Initialize the knot with an owner and public key. r.With(h.VerifySignature).Post("/init", h.Init) @@ -105,85 +106,3 @@ func Setup(ctx context.Context, c *config.Config, db *db.DB, e *rbac.Enforcer) ( return r, nil } - -func (h *Handle) StartJetstream(ctx context.Context) error { - collections := []string{tangled.PublicKeyNSID, tangled.KnotMemberNSID} - dids := []string{} - - var lastTimeUs int64 - var err error - lastTimeUs, err = h.db.GetLastTimeUs() - if err != nil { - log.Println("couldn't get last time us, starting from now") - lastTimeUs = time.Now().UnixMicro() - } - // If last time is older than a week, start from now - if time.Now().UnixMicro()-lastTimeUs > 7*24*60*60*1000*1000 { - lastTimeUs = time.Now().UnixMicro() - log.Printf("last time us is older than a week. discarding that and starting from now.") - err = h.db.SaveLastTimeUs(lastTimeUs) - if err != nil { - log.Println("failed to save last time us") - } - } - - log.Printf("found last time_us %d", lastTimeUs) - - h.js = jsclient.NewJetstreamClient(collections, dids) - messages, err := h.js.ReadJetstream(ctx, lastTimeUs) - if err != nil { - return fmt.Errorf("failed to read from jetstream: %w", err) - } - - go func() { - log.Println("waiting for knot to be initialized") - <-h.init - log.Println("initalized jetstream watcher") - - for msg := range messages { - var data map[string]interface{} - if err := json.Unmarshal(msg, &data); err != nil { - log.Printf("error unmarshaling message: %v", err) - continue - } - - if kind, ok := data["kind"].(string); ok && kind == "commit" { - commit := data["commit"].(map[string]interface{}) - - switch commit["collection"].(string) { - case tangled.PublicKeyNSID: - did := data["did"].(string) - record := commit["record"].(map[string]interface{}) - if err := h.db.AddPublicKeyFromRecord(did, record); err != nil { - log.Printf("failed to add public key: %v", err) - } else { - log.Printf("added public key from firehose: %s", data["did"]) - } - case tangled.KnotMemberNSID: - did := data["did"].(string) - record := commit["record"].(map[string]interface{}) - ok, err := h.e.E.Enforce(did, ThisServer, ThisServer, "server:invite") - if err != nil || !ok { - log.Printf("failed to add member from did %s", did) - } else { - log.Printf("adding member") - if err := h.e.AddMember(ThisServer, record["member"].(string)); err != nil { - log.Printf("failed to add member: %v", err) - } else { - log.Printf("added member from firehose: %s", record["member"]) - } - } - default: - } - - lastTimeUs := int64(data["time_us"].(float64)) - if err := h.db.SaveLastTimeUs(lastTimeUs); err != nil { - log.Printf("failed to save last time us: %v", err) - } - } - - } - }() - - return nil -} diff --git a/knotserver/jetstream.go b/knotserver/jetstream.go new file mode 100644 index 00000000..b16b4b67 --- /dev/null +++ b/knotserver/jetstream.go @@ -0,0 +1,144 @@ +package knotserver + +import ( + "context" + "encoding/json" + "fmt" + "io" + "log" + "net/http" + "path" + "strings" + "time" + + "github.com/sotangled/tangled/api/tangled" + "github.com/sotangled/tangled/knotserver/db" + "github.com/sotangled/tangled/knotserver/jsclient" +) + +func (h *Handle) StartJetstream(ctx context.Context) error { + collections := []string{tangled.PublicKeyNSID, tangled.KnotMemberNSID} + dids := []string{} + + lastTimeUs, err := h.getLastTimeUs() + if err != nil { + return err + } + + h.js = jsclient.NewJetstreamClient(collections, dids) + messages, err := h.js.ReadJetstream(ctx, lastTimeUs) + if err != nil { + return fmt.Errorf("failed to read from jetstream: %w", err) + } + + go h.processMessages(messages) + + return nil +} + +func (h *Handle) getLastTimeUs() (int64, error) { + lastTimeUs, err := h.db.GetLastTimeUs() + if err != nil { + log.Println("couldn't get last time us, starting from now") + lastTimeUs = time.Now().UnixMicro() + } + + // If last time is older than a week, start from now + if time.Now().UnixMicro()-lastTimeUs > 7*24*60*60*1000*1000 { + lastTimeUs = time.Now().UnixMicro() + log.Printf("last time us is older than a week. discarding that and starting from now.") + err = h.db.SaveLastTimeUs(lastTimeUs) + if err != nil { + log.Println("failed to save last time us") + } + } + + log.Printf("found last time_us %d", lastTimeUs) + return lastTimeUs, nil +} + +func (h *Handle) processPublicKey(did string, record map[string]interface{}) { + if err := h.db.AddPublicKeyFromRecord(did, record); err != nil { + log.Printf("failed to add public key: %v", err) + } else { + log.Printf("added public key from firehose: %s", did) + } +} + +func (h *Handle) fetchAndAddKeys(did string) { + resp, err := http.Get(path.Join(h.c.AppViewEndpoint, did)) + if err != nil { + log.Printf("error getting keys for %s: %v", did, err) + return + } + defer resp.Body.Close() + + plaintext, err := io.ReadAll(resp.Body) + if err != nil { + log.Printf("error reading response body: %v", err) + return + } + + for _, key := range strings.Split(string(plaintext), "\n") { + if key == "" { + continue + } + pk := db.PublicKey{ + Did: did, + } + pk.Key = key + if err := h.db.AddPublicKey(pk); err != nil { + log.Printf("failed to add public key: %v", err) + } + } +} + +func (h *Handle) processKnotMember(did string, record map[string]interface{}) { + ok, err := h.e.E.Enforce(did, ThisServer, ThisServer, "server:invite") + if err != nil || !ok { + log.Printf("failed to add member from did %s", did) + return + } + + log.Printf("adding member") + if err := h.e.AddMember(ThisServer, record["member"].(string)); err != nil { + log.Printf("failed to add member: %v", err) + } else { + log.Printf("added member from firehose: %s", record["member"]) + } + + h.fetchAndAddKeys(did) + h.js.UpdateDids([]string{did}) +} + +func (h *Handle) processMessages(messages <-chan []byte) { + log.Println("waiting for knot to be initialized") + <-h.init + log.Println("initalized jetstream watcher") + + for msg := range messages { + var data map[string]interface{} + if err := json.Unmarshal(msg, &data); err != nil { + log.Printf("error unmarshaling message: %v", err) + continue + } + + if kind, ok := data["kind"].(string); ok && kind == "commit" { + commit := data["commit"].(map[string]interface{}) + did := data["did"].(string) + record := commit["record"].(map[string]interface{}) + + switch commit["collection"].(string) { + case tangled.PublicKeyNSID: + h.processPublicKey(did, record) + case tangled.KnotMemberNSID: + h.processKnotMember(did, record) + } + + lastTimeUs := int64(data["time_us"].(float64)) + if err := h.db.SaveLastTimeUs(lastTimeUs); err != nil { + log.Printf("failed to save last time us: %v", err) + } + } + } +} diff --git a/knotserver/routes.go b/knotserver/routes.go index d4a78ee5..8f93af06 100644 --- a/knotserver/routes.go +++ b/knotserver/routes.go @@ -401,7 +401,40 @@ func (h *Handle) NewRepo(w http.ResponseWriter, r *http.Request) { w.WriteHeader(http.StatusNoContent) } -// TODO: make this set the initial user as the owner +func (h *Handle) AddMember(w http.ResponseWriter, r *http.Request) { + data := struct { + Did string `json:"did"` + PublicKeys []string `json:"keys"` + }{} + + if err := json.NewDecoder(r.Body).Decode(&data); err != nil { + writeError(w, "invalid request body", http.StatusBadRequest) + return + } + + did := data.Did + for _, k := range data.PublicKeys { + pk := db.PublicKey{ + Did: did, + } + pk.Key = k + err := h.db.AddPublicKey(pk) + if err != nil { + writeError(w, err.Error(), http.StatusInternalServerError) + return + } + } + + h.js.UpdateDids([]string{did}) + if err := h.e.AddMember(ThisServer, did); err != nil { + log.Println(err) + writeError(w, err.Error(), http.StatusInternalServerError) + return + } + + w.WriteHeader(http.StatusNoContent) +} + func (h *Handle) Init(w http.ResponseWriter, r *http.Request) { if h.knotInitialized { writeError(w, "knot already initialized", http.StatusConflict) @@ -441,7 +474,11 @@ func (h *Handle) Init(w http.ResponseWriter, r *http.Request) { } h.js.UpdateDids([]string{data.Did}) - h.e.AddOwner(ThisServer, data.Did) + if err := h.e.AddOwner(ThisServer, data.Did); err != nil { + log.Println(err) + writeError(w, err.Error(), http.StatusInternalServerError) + return + } // Signal that the knot is ready close(h.init)