1
0
Fork 0
mirror of https://github.com/sourcegraph/jsonrpc2.git synced 2026-06-16 04:04:56 +02:00

transparently simplify control flow (#83)

This commit is contained in:
Kevin Gillette 2025-02-17 07:55:54 -07:00 committed by GitHub
parent 534fd43609
commit 2cc94179e1
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
4 changed files with 72 additions and 89 deletions

61
conn.go
View file

@ -187,15 +187,17 @@ func (c *Conn) close(cause error) error {
} }
func (c *Conn) readMessages(ctx context.Context) { func (c *Conn) readMessages(ctx context.Context) {
var err error for {
for err == nil {
var m anyMessage var m anyMessage
err = c.stream.ReadObject(&m) err := c.stream.ReadObject(&m)
if err != nil { if err != nil {
break c.close(err)
return
} }
switch { switch {
// TODO: handle the case where both request and response are nil.
case m.request != nil: case m.request != nil:
for _, onRecv := range c.onRecv { for _, onRecv := range c.onRecv {
onRecv(m.request, nil) onRecv(m.request, nil)
@ -204,44 +206,37 @@ func (c *Conn) readMessages(ctx context.Context) {
case m.response != nil: case m.response != nil:
resp := m.response resp := m.response
if resp != nil {
id := resp.ID id := resp.ID
c.mu.Lock() c.mu.Lock()
call := c.pending[id] call := c.pending[id]
delete(c.pending, id) delete(c.pending, id)
c.mu.Unlock() c.mu.Unlock()
if call != nil {
call.response = resp
}
if len(c.onRecv) > 0 {
var req *Request var req *Request
if call != nil { if call != nil {
call.response = resp
req = call.request req = call.request
} }
for _, onRecv := range c.onRecv { for _, onRecv := range c.onRecv {
onRecv(req, resp) onRecv(req, resp)
} }
}
switch { if call == nil {
case call == nil:
c.logger.Printf("jsonrpc2: ignoring response #%s with no corresponding request\n", id) c.logger.Printf("jsonrpc2: ignoring response #%s with no corresponding request\n", id)
continue
}
case resp.Error != nil: var err error
call.done <- resp.Error if resp.Error != nil {
close(call.done) err = resp.Error
}
default: call.done <- err
call.done <- nil
close(call.done) close(call.done)
} }
} }
} }
}
c.close(err)
}
func (c *Conn) send(_ context.Context, m *anyMessage, wait bool) (cc *call, err error) { func (c *Conn) send(_ context.Context, m *anyMessage, wait bool) (cc *call, err error) {
c.sending.Lock() c.sending.Lock()
@ -339,25 +334,20 @@ type Waiter struct {
// error is returned. // error is returned.
func (w Waiter) Wait(ctx context.Context, result interface{}) error { func (w Waiter) Wait(ctx context.Context, result interface{}) error {
select { select {
case <-ctx.Done():
return ctx.Err()
case err, ok := <-w.call.done: case err, ok := <-w.call.done:
if !ok { if !ok {
err = ErrClosed return ErrClosed
} }
if err != nil { if err != nil || result == nil {
return err return err
} }
if result != nil {
if w.call.response.Result == nil { if w.call.response.Result == nil {
w.call.response.Result = &jsonNull w.call.response.Result = &jsonNull
} }
if err := json.Unmarshal(*w.call.response.Result, result); err != nil { return json.Unmarshal(*w.call.response.Result, result)
return err
}
}
return nil
case <-ctx.Done():
return ctx.Err()
} }
} }
@ -423,12 +413,7 @@ func (m *anyMessage) UnmarshalJSON(data []byte) error {
return errors.New("jsonrpc2: invalid empty batch") return errors.New("jsonrpc2: invalid empty batch")
} }
for i := range msgs { for i := range msgs {
if err := checkType(&msg{ if err := checkType(&msgs[i]); err != nil {
ID: msgs[i].ID,
Method: msgs[i].Method,
Result: msgs[i].Result,
Error: msgs[i].Error,
}); err != nil {
return err return err
} }
} }

View file

@ -44,11 +44,9 @@ func LogMessages(logger Logger) ConnOpt {
OnRecv(func(req *Request, resp *Response) { OnRecv(func(req *Request, resp *Response) {
switch { switch {
case resp != nil: case resp != nil:
var method string method := "(no matching request)"
if req != nil { if req != nil {
method = req.Method method = req.Method
} else {
method = "(no matching request)"
} }
switch { switch {
case resp.Result != nil: case resp.Result != nil:

View file

@ -30,22 +30,18 @@ func (h *HandlerWithErrorConfigurer) Handle(ctx context.Context, conn *Conn, req
if err == nil { if err == nil {
err = resp.SetResult(result) err = resp.SetResult(result)
} }
if err != nil {
if e, ok := err.(*Error); ok { if e, ok := err.(*Error); ok {
resp.Error = e resp.Error = e
} else { } else if err != nil {
resp.Error = &Error{Message: err.Error()} resp.Error = &Error{Message: err.Error()}
} }
}
if !req.Notif { err = conn.SendResponse(ctx, resp)
if err := conn.SendResponse(ctx, resp); err != nil { if err != nil && (err != ErrClosed || !h.suppressErrClosed) {
if err != ErrClosed || !h.suppressErrClosed {
conn.logger.Printf("jsonrpc2 handler: sending response %s: %v\n", resp.ID, err) conn.logger.Printf("jsonrpc2 handler: sending response %s: %v\n", resp.ID, err)
} }
} }
}
}
// SuppressErrClosed makes the handler suppress jsonrpc2.ErrClosed errors from // SuppressErrClosed makes the handler suppress jsonrpc2.ErrClosed errors from
// being logged. The original handler `h` is returned. // being logged. The original handler `h` is returned.

View file

@ -55,6 +55,10 @@ func (r Request) MarshalJSON() ([]byte, error) {
// UnmarshalJSON implements json.Unmarshaler. // UnmarshalJSON implements json.Unmarshaler.
func (r *Request) UnmarshalJSON(data []byte) error { func (r *Request) UnmarshalJSON(data []byte) error {
r2 := make(map[string]interface{}) r2 := make(map[string]interface{})
pop := func(key string) interface{} {
defer delete(r2, key)
return r2[key]
}
// Detect if the "params" or "meta" fields are JSON "null" or just not // Detect if the "params" or "meta" fields are JSON "null" or just not
// present by seeing if the field gets overwritten to nil. // present by seeing if the field gets overwritten to nil.
@ -68,36 +72,37 @@ func (r *Request) UnmarshalJSON(data []byte) error {
if err := decoder.Decode(&r2); err != nil { if err := decoder.Decode(&r2); err != nil {
return err return err
} }
var ok bool var ok bool
r.Method, ok = r2["method"].(string) r.Method, ok = pop("method").(string)
if !ok { if !ok {
return errors.New("missing method field") return errors.New("missing method field")
} }
switch { switch params := pop("params"); params {
case r2["params"] == nil: case nil:
r.Params = &jsonNull r.Params = &jsonNull
case r2["params"] == emptyParams: case emptyParams:
r.Params = nil r.Params = nil
default: default:
b, err := json.Marshal(r2["params"]) b, err := json.Marshal(params)
if err != nil { if err != nil {
return fmt.Errorf("failed to marshal params: %w", err) return fmt.Errorf("failed to marshal params: %w", err)
} }
r.Params = (*json.RawMessage)(&b) r.Params = (*json.RawMessage)(&b)
} }
switch { switch meta := pop("meta"); meta {
case r2["meta"] == nil: case nil:
r.Meta = &jsonNull r.Meta = &jsonNull
case r2["meta"] == emptyMeta: case emptyMeta:
r.Meta = nil r.Meta = nil
default: default:
b, err := json.Marshal(r2["meta"]) b, err := json.Marshal(meta)
if err != nil { if err != nil {
return fmt.Errorf("failed to marshal Meta: %w", err) return fmt.Errorf("failed to marshal Meta: %w", err)
} }
r.Meta = (*json.RawMessage)(&b) r.Meta = (*json.RawMessage)(&b)
} }
switch rawID := r2["id"].(type) { switch rawID := pop("id").(type) {
case nil: case nil:
r.ID = ID{} r.ID = ID{}
r.Notif = true r.Notif = true
@ -115,13 +120,12 @@ func (r *Request) UnmarshalJSON(data []byte) error {
return fmt.Errorf("unexpected ID type: %T", rawID) return fmt.Errorf("unexpected ID type: %T", rawID)
} }
// The jsonrpc field should not be added to ExtraFields.
delete(r2, "jsonrpc")
// Clear the extra fields before populating them again. // Clear the extra fields before populating them again.
r.ExtraFields = nil r.ExtraFields = nil
for name, value := range r2 { for name, value := range r2 {
switch name {
case "id", "jsonrpc", "meta", "method", "params":
continue
}
r.ExtraFields = append(r.ExtraFields, RequestField{ r.ExtraFields = append(r.ExtraFields, RequestField{
Name: name, Name: name,
Value: value, Value: value,