Skip to content
Closed
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
13 changes: 10 additions & 3 deletions core/core.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,21 +6,24 @@ import (
"sync"
)

// KeyValueStore is a thread-safe in-memory key-value store with an optional transaction logger.
type KeyValueStore struct {
sync.RWMutex
m map[string]string
transact TransactionLogger
}

var ErrorNoSuchKey = errors.New("no suck key")
var ErrorNoSuchKey = errors.New("no such key")

// NewKeyValueStore initializes and returns a new KeyValueStore.
func NewKeyValueStore() *KeyValueStore {
return &KeyValueStore{
m: make(map[string]string),
transact: ZeroTransactionLogger{},
}
}

// Delete removes a key from the store and logs the deletion.
func (store *KeyValueStore) Delete(key string) error {
store.Lock()
delete(store.m, key)
Expand All @@ -31,6 +34,7 @@ func (store *KeyValueStore) Delete(key string) error {
return nil
}

// Put inserts or updates a key-value pair in the store and logs the operation.
func (store *KeyValueStore) Put(key, value string) error {
store.Lock()
store.m[key] = value
Expand All @@ -41,6 +45,7 @@ func (store *KeyValueStore) Put(key, value string) error {
return nil
}

// Get retrieves the value associated with a key from the store.
func (store *KeyValueStore) Get(key string) (string, error) {
store.RLock()
value, ok := store.m[key]
Expand All @@ -53,20 +58,22 @@ func (store *KeyValueStore) Get(key string) (string, error) {
return value, nil
}

// WithTransactionLogger sets the transaction logger for the store and returns the store.
func (store *KeyValueStore) WithTransactionLogger(tl TransactionLogger) *KeyValueStore {
store.transact = tl
return store
}

// Restore replays events from the transaction logger to rebuild the store's state.
func (store *KeyValueStore) Restore() error {
var err error

events, errors := store.transact.ReadEvents()
events, errs := store.transact.ReadEvents()
count, ok, e := 0, true, Event{}

for ok && err == nil {
select {
case err, ok = <-errors:
case err, ok = <-errs:

case e, ok = <-events:
switch e.EventType {
Expand Down
6 changes: 3 additions & 3 deletions frontend/grpc.go
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,7 @@ func (g *grpcFrontend) Get(ctx context.Context, r *pb.GetRequest) (*pb.GetRespon
}

func (g *grpcFrontend) Put(ctx context.Context, r *pb.PutRequest) (*pb.PutResponse, error) {
log.Printf("Received PUT key=%v value%v\n", r.Key, r.Value)
log.Printf("Received PUT key=%v value=%v\n", r.Key, r.Value)

err := g.store.Put(r.Key, string(r.Value))

Expand All @@ -49,12 +49,12 @@ func (g *grpcFrontend) Start(store *core.KeyValueStore) error {

lis, err := net.Listen("tcp", ":50051")
if err != nil {
return fmt.Errorf("Failed to listen: %v", err)
return fmt.Errorf("failed to listen: %v", err)
}

fmt.Println("Listening on :50051")
if err := s.Serve(lis); err != nil {
return fmt.Errorf("failed to server: %v", err)
return fmt.Errorf("failed to serve: %v", err)
}

return nil
Expand Down
8 changes: 4 additions & 4 deletions frontend/rest.go
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,7 @@

r := mux.NewRouter()

r.Use(f.logginMiddleware)
r.Use(f.loggingMiddleware)

r.HandleFunc("/v1/{key}", f.keyValueGetHandler).Methods("GET")
r.HandleFunc("/v1/{key}", f.keyValuePutHandler).Methods("PUT")
Expand All @@ -40,7 +40,7 @@
return http.ListenAndServe(":"+port, r)
}

func (f *restFrontend) logginMiddleware(next http.Handler) http.Handler {
func (f *restFrontend) loggingMiddleware(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
log.Println(r.Method, r.RequestURI)
next.ServeHTTP(w, r)
Expand Down Expand Up @@ -80,12 +80,12 @@
return
}

w.Write([]byte(value))

Check failure on line 83 in frontend/rest.go

View workflow job for this annotation

GitHub Actions / Lint

Error return value of `w.Write` is not checked (errcheck)
}

func (f *restFrontend) keyValuePutHandler(w http.ResponseWriter, r *http.Request) {
vars := mux.Vars(r)
keys := vars["key"]
key := vars["key"]
value, err := io.ReadAll(r.Body)

defer r.Body.Close()
Expand All @@ -96,7 +96,7 @@
http.StatusInternalServerError)
return
}
err = f.store.Put(keys, string(value))
err = f.store.Put(key, string(value))
if err != nil {
http.Error(w,
err.Error(),
Expand Down
4 changes: 1 addition & 3 deletions transact/filelogger.go
Original file line number Diff line number Diff line change
Expand Up @@ -57,8 +57,6 @@ func (l *FileTransactionLogger) Run() {
l.errors <- err
return
}

l.wg.Wait()
}
}()
}
Expand Down Expand Up @@ -116,7 +114,7 @@ func (l *FileTransactionLogger) ReadEvents() (<-chan core.Event, <-chan error) {
}

if err := scanner.Err(); err != nil {
outError <- fmt.Errorf("Transaction log failed to reader : %w", err)
outError <- fmt.Errorf("transaction log read failure: %w", err)
return
}
}()
Expand Down
5 changes: 4 additions & 1 deletion transact/pglogger.go
Original file line number Diff line number Diff line change
Expand Up @@ -153,7 +153,10 @@ func NewPostgresTransactionLogger(args PostgresDBParams) (core.TransactionLogger
return nil, fmt.Errorf("failed to open db connection: %w", err)
}

logger := &PostgresTransactionLogger{db: db}
logger := &PostgresTransactionLogger{
db: db,
wg: &sync.WaitGroup{},
}

logger.CreateTable()

Expand Down
Loading