Docs / ArchitectureKit / Guides / Putting It Together

Putting It Together

So far, the pieces have been shown one at a time. A main function wires them together: it creates the store, registers the schemas, starts the projection and waits until it has caught up, serves the routes and the health checks, and shuts down cleanly once it is asked to stop:

func main() {
  ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
  defer stop()

  baseURL, err := url.Parse("http://localhost:3000")
  if err != nil {
    log.Fatal(err)
  }

  client, err := eventsourcingdb.NewClient(baseURL, "secret")
  if err != nil {
    log.Fatal(err)
  }

  store := architecturekit.NewStore(client, "https://library.eventsourcingdb.io")

  if err := architecturekit.RegisterSchemas(ctx, store, bookState.Schemas()); err != nil {
    log.Fatal(err)
  }

  // The projection stops only after the server, since the requests that the
  // server finishes may still wait for the view.
  projectionCtx, stopProjection := context.WithCancel(context.Background())
  defer stopProjection()

  catalog := newCatalog()
  run := architecturekit.StartProjection(projectionCtx, store, architecturekit.SubjectTree("/books"),
    architecturekit.Tracking(newCatalogProjection(catalog), catalog),
    architecturekit.Named("catalog"),
  )

  select {
  case <-run.CaughtUp():
  case <-run.Done():
    log.Fatal(run.Err())
  case <-time.After(time.Minute):
    log.Fatal("the catalog has not caught up within a minute")
  case <-ctx.Done():
    return
  }

  api := httpapi.NewAPI(store, userFrom)
  mux := http.NewServeMux()

  httpapi.Route(api, mux, "POST /api/books/{id}/borrow", toBorrowBook, borrowBook)
  httpapi.Query(api, mux, "GET /api/books", toListBooks, answerBooks(listBooks(catalog)),
    httpapi.Revisioned(catalog, httpapi.DefaultWait),
  )
  mux.Handle("GET /ready", httpapi.Readiness(run))
  mux.Handle("GET /live", httpapi.Liveness(run))

  server := &http.Server{Addr: ":8080", Handler: mux}
  go func() {
    if err := server.ListenAndServe(); !errors.Is(err, http.ErrServerClosed) {
      log.Fatal(err)
    }
  }()

  <-ctx.Done()

  shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
  defer cancel()

  if err := server.Shutdown(shutdownCtx); err != nil {
    log.Println(err)
  }

  stopProjection()
  <-run.Done()

  if err := run.Err(); err != nil {
    log.Println(err)
  }
}

signal.NotifyContext ends ctx once the process is asked to stop, as an orchestrator does with SIGTERM. Shutdown then lets the server finish the requests it has begun, which may still wait for the view (see Reading Your Own Writes over HTTP). That is why the projection runs with a context of its own, which main cancels only after the server has stopped. Then it waits for Done, and logs the error of the run, which is nil unless the run had ended on a failure before.

A database that can not be reached at the start makes RegisterSchemas fail right away. The server starts only once the view has caught up, so that no query sees a half-built view. If that takes longer than a minute, for example because the database has stopped answering in the meantime, main gives up rather than wait without end, so that the orchestrator restarts the application. Choose the limit with room for the history to grow, since a view that never catches up within it keeps the application from ever starting. Until the server starts, the health checks do not answer either, so give the application that long to start, for example with a startup probe in Kubernetes.

Set up this way, everything lives as long as main does. If you move the setup of the projection into a function of its own, that function returns long before the application ends, so it must not cancel the context of the projection with defer cancel(), which would stop the projection right away. Have it return the cancel function along with the run instead, and call it on shutdown, as main calls stopProjection.