Skip to content
Open
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
12 changes: 10 additions & 2 deletions bulker/config-keeper/app.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ package main

import (
"context"
"encoding/json"
"fmt"
"github.com/jitsucom/bulker/jitsubase/appbase"
"io"
Expand All @@ -19,14 +20,21 @@ type Context struct {
}

type RawRepositoryData struct {
data atomic.Pointer[[]byte]
// validateJSON rejects payloads that are not complete, valid JSON. A repository
// source that fails mid-stream can produce a truncated body with HTTP 200 —
// without this check such payload gets cached and served to every consumer.
validateJSON bool
data atomic.Pointer[[]byte]
}

func (r *RawRepositoryData) Init(reader io.Reader, tag any) error {
data, err := io.ReadAll(reader)
if err != nil {
return err
}
if r.validateJSON && !json.Valid(data) {
return fmt.Errorf("payload is not valid JSON (%d bytes) - keeping previous data", len(data))
}
r.data.Store(&data)
return nil
}
Expand Down Expand Up @@ -59,7 +67,7 @@ func (a *Context) InitContext(settings *appbase.AppSettings) error {
"p.js": a.pScript,
}
for _, rep := range strings.Split(reps, ",") {
a.repositories[rep] = appbase.NewHTTPRepository[[]byte](rep, baseUrl+"/"+rep, token, appbase.HTTPTagLastModified, &RawRepositoryData{}, 2, refreshPeriodSec, cacheDir)
a.repositories[rep] = appbase.NewHTTPRepository[[]byte](rep, baseUrl+"/"+rep, token, appbase.HTTPTagLastModified, &RawRepositoryData{validateJSON: true}, 2, refreshPeriodSec, cacheDir)

}
router := NewRouter(a)
Expand Down
2 changes: 1 addition & 1 deletion bulker/config-keeper/router.go
Original file line number Diff line number Diff line change
Expand Up @@ -77,7 +77,7 @@ func (r *Router) RepositoryHandler(c *gin.Context) {
repository, ok := r.appContext.repositories[repName]
if !ok {
r.Infof("Repository %s not found, initializing", repName)
repository = appbase.NewHTTPRepository[[]byte](repName, r.appContext.config.RepositoryBaseURL+"/"+repName, r.appContext.config.RepositoryAuthToken, appbase.HTTPTagLastModified, &RawRepositoryData{}, 2, r.appContext.config.RepositoryRefreshPeriodSec, r.appContext.config.CacheDir)
repository = appbase.NewHTTPRepository[[]byte](repName, r.appContext.config.RepositoryBaseURL+"/"+repName, r.appContext.config.RepositoryAuthToken, appbase.HTTPTagLastModified, &RawRepositoryData{validateJSON: true}, 2, r.appContext.config.RepositoryRefreshPeriodSec, r.appContext.config.CacheDir)
initTimeout := time.After(time.Second * 60)
ticker := time.NewTicker(time.Second)
defer ticker.Stop()
Expand Down
8 changes: 6 additions & 2 deletions bulker/jitsubase/logging/global_logger.go
Original file line number Diff line number Diff line change
Expand Up @@ -121,10 +121,14 @@ func Warn(v ...any) {
log.Warnln(v...)
}

// Fatal-level failures abort the process (typically failure to start), so they
// carry the "System error:" marker used by log-based alerting.
func Fatal(v ...any) {
log.Fatal(v...)
msg := []any{"System error:"}
msg = append(msg, v...)
log.Fatal(msg...)
}

func Fatalf(format string, v ...any) {
log.Fatalf(format, v...)
log.Fatalf("System error: "+format, v...)
}
12 changes: 11 additions & 1 deletion libs/core-functions-lib/src/lib/inmem-store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -75,7 +75,10 @@ export const createInMemoryStore = <T>(definition: StoreDefinition<T>): InMemory
status = "ok";
lastRefresh = new Date();
} catch (e) {
log.atWarn().withCause(e).log(`Failed to refresh store ${definition.name}. Using an old value`);
// Not a system error (the store keeps serving the old value) — but the message
// wording matches the Go-side repository refresh error in bulker/jitsubase so
// one log query covers both stacks
log.atError().withCause(e).log(`Error refreshing repository ${definition.name}. Using an old value`);
status = "outdated";
}
};
Expand Down Expand Up @@ -109,6 +112,13 @@ export const createInMemoryStore = <T>(definition: StoreDefinition<T>): InMemory
const cachedInstance = loadFromCache(definition);
if (!cachedInstance) {
status = "failed";
log
.atError()
.log(
`System error: Failed to initialize store ${definition.name}. Initial load failed with ${getErrorMessage(
e
)} and no local cache found`
);
reject(
new Error(
`Failed to initialize store ${definition.name}. Initial load failed with ${getErrorMessage(
Expand Down
Loading
Loading