Expose live query error messages via API (#205)

Somewhere around osquery 4.4.0 these messages were added to query
responses. We can now expose them to the API clients rather than using
the placeholder text.

Required for #192
This commit is contained in:
Zach Wasserman
2021-01-19 14:52:29 -08:00
committed by GitHub
parent f254a9a343
commit deaf8880f3
10 changed files with 64 additions and 42 deletions
+1 -1
View File
@@ -17,7 +17,7 @@ type OsqueryService interface {
// for) should be returned. Returning 0 for this will not activate the
// feature.
GetDistributedQueries(ctx context.Context) (queries map[string]string, accelerate uint, err error)
SubmitDistributedQueryResults(ctx context.Context, results OsqueryDistributedQueryResults, statuses map[string]OsqueryStatus) (err error)
SubmitDistributedQueryResults(ctx context.Context, results OsqueryDistributedQueryResults, statuses map[string]OsqueryStatus, messages map[string]string) (err error)
SubmitStatusLogs(ctx context.Context, logs []json.RawMessage) (err error)
SubmitResultLogs(ctx context.Context, logs []json.RawMessage) (err error)
//CarveBegin(ctx context.Context)
+3 -1
View File
@@ -120,7 +120,9 @@ func (svc *launcherWrapper) PublishResults(ctx context.Context, nodeKey string,
osqueryResults[result.QueryName] = result.Rows
}
err = svc.tls.SubmitDistributedQueryResults(newCtx, osqueryResults, statuses)
// TODO can Launcher expose the error messages?
messages := make(map[string]string)
err = svc.tls.SubmitDistributedQueryResults(newCtx, osqueryResults, statuses, messages)
return "", "", false, errors.Wrap(err, "submit launcher results")
}
+5 -2
View File
@@ -5,10 +5,10 @@ import (
"encoding/json"
"testing"
"github.com/go-kit/kit/log"
"github.com/fleetdm/fleet/server/health"
"github.com/fleetdm/fleet/server/kolide"
"github.com/fleetdm/fleet/server/mock"
"github.com/go-kit/kit/log"
"github.com/kolide/osquery-go/plugin/distributed"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
@@ -66,7 +66,9 @@ func TestLauncherPublishResults(t *testing.T) {
tls.SubmitDistributedQueryResultsFunc = func(
ctx context.Context,
results kolide.OsqueryDistributedQueryResults,
statuses map[string]kolide.OsqueryStatus) (err error) {
statuses map[string]kolide.OsqueryStatus,
messages map[string]string,
) (err error) {
assert.Equal(t, results["query"][0], result)
return nil
}
@@ -146,6 +148,7 @@ func newTLSService(t *testing.T) *mock.TLSService {
ctx context.Context,
results kolide.OsqueryDistributedQueryResults,
statuses map[string]kolide.OsqueryStatus,
messages map[string]string,
) (err error) {
return
},
+3 -3
View File
@@ -19,7 +19,7 @@ type GetClientConfigFunc func(ctx context.Context) (config map[string]interface{
type GetDistributedQueriesFunc func(ctx context.Context) (queries map[string]string, accelerate uint, err error)
type SubmitDistributedQueryResultsFunc func(ctx context.Context, results kolide.OsqueryDistributedQueryResults, statuses map[string]kolide.OsqueryStatus) (err error)
type SubmitDistributedQueryResultsFunc func(ctx context.Context, results kolide.OsqueryDistributedQueryResults, statuses map[string]kolide.OsqueryStatus, messages map[string]string) (err error)
type SubmitStatusLogsFunc func(ctx context.Context, logs []json.RawMessage) (err error)
@@ -68,9 +68,9 @@ func (s *TLSService) GetDistributedQueries(ctx context.Context) (queries map[str
return s.GetDistributedQueriesFunc(ctx)
}
func (s *TLSService) SubmitDistributedQueryResults(ctx context.Context, results kolide.OsqueryDistributedQueryResults, statuses map[string]kolide.OsqueryStatus) (err error) {
func (s *TLSService) SubmitDistributedQueryResults(ctx context.Context, results kolide.OsqueryDistributedQueryResults, statuses map[string]kolide.OsqueryStatus, messages map[string]string) (err error) {
s.SubmitDistributedQueryResultsFuncInvoked = true
return s.SubmitDistributedQueryResultsFunc(ctx, results, statuses)
return s.SubmitDistributedQueryResultsFunc(ctx, results, statuses, messages)
}
func (s *TLSService) SubmitStatusLogs(ctx context.Context, logs []json.RawMessage) (err error) {
+3 -2
View File
@@ -4,8 +4,8 @@ import (
"context"
"encoding/json"
"github.com/go-kit/kit/endpoint"
"github.com/fleetdm/fleet/server/kolide"
"github.com/go-kit/kit/endpoint"
)
////////////////////////////////////////////////////////////////////////////////
@@ -98,6 +98,7 @@ type submitDistributedQueryResultsRequest struct {
NodeKey string `json:"node_key"`
Results kolide.OsqueryDistributedQueryResults `json:"queries"`
Statuses map[string]kolide.OsqueryStatus `json:"statuses"`
Messages map[string]string `json:"messages"`
}
type submitDistributedQueryResultsResponse struct {
@@ -109,7 +110,7 @@ func (r submitDistributedQueryResultsResponse) error() error { return r.Err }
func makeSubmitDistributedQueryResultsEndpoint(svc kolide.Service) endpoint.Endpoint {
return func(ctx context.Context, request interface{}) (interface{}, error) {
req := request.(submitDistributedQueryResultsRequest)
err := svc.SubmitDistributedQueryResults(ctx, req.Results, req.Statuses)
err := svc.SubmitDistributedQueryResults(ctx, req.Results, req.Statuses, req.Messages)
if err != nil {
return submitDistributedQueryResultsResponse{Err: err}, nil
}
+2 -2
View File
@@ -92,7 +92,7 @@ func (mw loggingMiddleware) GetDistributedQueries(ctx context.Context) (map[stri
return queries, accelerate, err
}
func (mw loggingMiddleware) SubmitDistributedQueryResults(ctx context.Context, results kolide.OsqueryDistributedQueryResults, statuses map[string]kolide.OsqueryStatus) error {
func (mw loggingMiddleware) SubmitDistributedQueryResults(ctx context.Context, results kolide.OsqueryDistributedQueryResults, statuses map[string]kolide.OsqueryStatus, messages map[string]string) error {
var (
err error
)
@@ -107,7 +107,7 @@ func (mw loggingMiddleware) SubmitDistributedQueryResults(ctx context.Context, r
)
}(time.Now())
err = mw.Service.SubmitDistributedQueryResults(ctx, results, statuses)
err = mw.Service.SubmitDistributedQueryResults(ctx, results, statuses, messages)
return err
}
+5 -8
View File
@@ -9,10 +9,10 @@ import (
"strings"
"time"
"github.com/go-kit/kit/log"
hostctx "github.com/fleetdm/fleet/server/contexts/host"
"github.com/fleetdm/fleet/server/kolide"
"github.com/fleetdm/fleet/server/pubsub"
"github.com/go-kit/kit/log"
"github.com/pkg/errors"
"github.com/spf13/cast"
)
@@ -567,7 +567,7 @@ func (svc service) ingestLabelQuery(host kolide.Host, query string, rows []map[s
// ingestDistributedQuery takes the results of a distributed query and modifies the
// provided kolide.Host appropriately.
func (svc service) ingestDistributedQuery(host kolide.Host, name string, rows []map[string]string, failed bool) error {
func (svc service) ingestDistributedQuery(host kolide.Host, name string, rows []map[string]string, failed bool, errMsg string) error {
trimmedQuery := strings.TrimPrefix(name, hostDistributedQueryPrefix)
campaignID, err := strconv.Atoi(emptyToZero(trimmedQuery))
@@ -582,10 +582,7 @@ func (svc service) ingestDistributedQuery(host kolide.Host, name string, rows []
Rows: rows,
}
if failed {
// osquery errors are not currently helpful, but we should fix
// them to be better in the future
errString := "failed"
res.Error = &errString
res.Error = &errMsg
}
err = svc.resultStore.WriteResult(res)
@@ -632,7 +629,7 @@ func (svc service) ingestDistributedQuery(host kolide.Host, name string, rows []
return nil
}
func (svc service) SubmitDistributedQueryResults(ctx context.Context, results kolide.OsqueryDistributedQueryResults, statuses map[string]kolide.OsqueryStatus) error {
func (svc service) SubmitDistributedQueryResults(ctx context.Context, results kolide.OsqueryDistributedQueryResults, statuses map[string]kolide.OsqueryStatus, messages map[string]string) error {
host, ok := hostctx.FromContext(ctx)
if !ok {
@@ -659,7 +656,7 @@ func (svc service) SubmitDistributedQueryResults(ctx context.Context, results ko
// status indicates a query error
status, ok := statuses[query]
failed := (ok && status != kolide.StatusOK)
err = svc.ingestDistributedQuery(host, query, rows, failed)
err = svc.ingestDistributedQuery(host, query, rows, failed, messages[query])
default:
err = osqueryError{message: "unknown query prefix: " + query}
}
+14 -12
View File
@@ -12,7 +12,6 @@ import (
"time"
"github.com/WatchBeam/clock"
"github.com/go-kit/kit/log"
"github.com/fleetdm/fleet/server/config"
hostctx "github.com/fleetdm/fleet/server/contexts/host"
"github.com/fleetdm/fleet/server/contexts/viewer"
@@ -22,6 +21,7 @@ import (
"github.com/fleetdm/fleet/server/logging"
"github.com/fleetdm/fleet/server/mock"
"github.com/fleetdm/fleet/server/pubsub"
"github.com/go-kit/kit/log"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)
@@ -347,6 +347,7 @@ func TestLabelQueries(t *testing.T) {
hostLabelQueryPrefix + "1": {{"col1": "val1"}},
},
map[string]kolide.OsqueryStatus{},
map[string]string{},
)
assert.Nil(t, err)
host.LabelUpdateTime = mockClock.Now()
@@ -366,6 +367,7 @@ func TestLabelQueries(t *testing.T) {
hostLabelQueryPrefix + "3": {},
},
map[string]kolide.OsqueryStatus{},
map[string]string{},
)
assert.Nil(t, err)
host.LabelUpdateTime = mockClock.Now()
@@ -643,7 +645,7 @@ func TestDetailQueriesWithEmptyStrings(t *testing.T) {
}
// Verify that results are ingested properly
svc.SubmitDistributedQueryResults(ctx, results, map[string]kolide.OsqueryStatus{})
svc.SubmitDistributedQueryResults(ctx, results, map[string]kolide.OsqueryStatus{}, map[string]string{})
// osquery_info
assert.Equal(t, "darwin", gotHost.Platform)
@@ -813,7 +815,7 @@ func TestDetailQueries(t *testing.T) {
return nil
}
// Verify that results are ingested properly
svc.SubmitDistributedQueryResults(ctx, results, map[string]kolide.OsqueryStatus{})
svc.SubmitDistributedQueryResults(ctx, results, map[string]kolide.OsqueryStatus{}, map[string]string{})
// osquery_info
assert.Equal(t, "darwin", gotHost.Platform)
@@ -1069,7 +1071,7 @@ func TestDistributedQueryResults(t *testing.T) {
// this test.
time.Sleep(10 * time.Millisecond)
err = svc.SubmitDistributedQueryResults(hostCtx, results, map[string]kolide.OsqueryStatus{})
err = svc.SubmitDistributedQueryResults(hostCtx, results, map[string]kolide.OsqueryStatus{}, map[string]string{})
require.Nil(t, err)
}
@@ -1087,7 +1089,7 @@ func TestIngestDistributedQueryParseIdError(t *testing.T) {
}
host := kolide.Host{ID: 1}
err := svc.ingestDistributedQuery(host, "bad_name", []map[string]string{}, false)
err := svc.ingestDistributedQuery(host, "bad_name", []map[string]string{}, false, "")
require.Error(t, err)
assert.Contains(t, err.Error(), "unable to parse campaign")
}
@@ -1111,7 +1113,7 @@ func TestIngestDistributedQueryOrphanedCampaignLoadError(t *testing.T) {
host := kolide.Host{ID: 1}
err := svc.ingestDistributedQuery(host, "kolide_distributed_query_42", []map[string]string{}, false)
err := svc.ingestDistributedQuery(host, "kolide_distributed_query_42", []map[string]string{}, false, "")
require.Error(t, err)
assert.Contains(t, err.Error(), "loading orphaned campaign")
}
@@ -1144,7 +1146,7 @@ func TestIngestDistributedQueryOrphanedCampaignWaitListener(t *testing.T) {
host := kolide.Host{ID: 1}
err := svc.ingestDistributedQuery(host, "kolide_distributed_query_42", []map[string]string{}, false)
err := svc.ingestDistributedQuery(host, "kolide_distributed_query_42", []map[string]string{}, false, "")
require.Error(t, err)
assert.Contains(t, err.Error(), "campaign waiting for listener")
}
@@ -1180,7 +1182,7 @@ func TestIngestDistributedQueryOrphanedCloseError(t *testing.T) {
host := kolide.Host{ID: 1}
err := svc.ingestDistributedQuery(host, "kolide_distributed_query_42", []map[string]string{}, false)
err := svc.ingestDistributedQuery(host, "kolide_distributed_query_42", []map[string]string{}, false, "")
require.Error(t, err)
assert.Contains(t, err.Error(), "closing orphaned campaign")
}
@@ -1217,7 +1219,7 @@ func TestIngestDistributedQueryOrphanedStopError(t *testing.T) {
host := kolide.Host{ID: 1}
err := svc.ingestDistributedQuery(host, "kolide_distributed_query_42", []map[string]string{}, false)
err := svc.ingestDistributedQuery(host, "kolide_distributed_query_42", []map[string]string{}, false, "")
require.Error(t, err)
assert.Contains(t, err.Error(), "stopping orphaned campaign")
}
@@ -1254,7 +1256,7 @@ func TestIngestDistributedQueryOrphanedStop(t *testing.T) {
host := kolide.Host{ID: 1}
err := svc.ingestDistributedQuery(host, "kolide_distributed_query_42", []map[string]string{}, false)
err := svc.ingestDistributedQuery(host, "kolide_distributed_query_42", []map[string]string{}, false, "")
require.NoError(t, err)
lq.AssertExpectations(t)
}
@@ -1284,7 +1286,7 @@ func TestIngestDistributedQueryRecordCompletionError(t *testing.T) {
}()
time.Sleep(10 * time.Millisecond)
err := svc.ingestDistributedQuery(host, "kolide_distributed_query_42", []map[string]string{}, false)
err := svc.ingestDistributedQuery(host, "kolide_distributed_query_42", []map[string]string{}, false, "")
require.Error(t, err)
assert.Contains(t, err.Error(), "record query completion")
lq.AssertExpectations(t)
@@ -1315,7 +1317,7 @@ func TestIngestDistributedQuery(t *testing.T) {
}()
time.Sleep(10 * time.Millisecond)
err := svc.ingestDistributedQuery(host, "kolide_distributed_query_42", []map[string]string{}, false)
err := svc.ingestDistributedQuery(host, "kolide_distributed_query_42", []map[string]string{}, false, "")
require.NoError(t, err)
lq.AssertExpectations(t)
}
+3
View File
@@ -52,10 +52,12 @@ func decodeSubmitDistributedQueryResultsRequest(ctx context.Context, r *http.Req
// },
// "node_key":"IGXCXknWQ1baTa8TZ6rF3kAPZ4\/aTsui"
// }
type distributedQueryResultsShim struct {
NodeKey string `json:"node_key"`
Results map[string]json.RawMessage `json:"queries"`
Statuses map[string]interface{} `json:"statuses"`
Messages map[string]string `json:"messages"`
}
var shim distributedQueryResultsShim
@@ -97,6 +99,7 @@ func decodeSubmitDistributedQueryResultsRequest(ctx context.Context, r *http.Req
NodeKey: shim.NodeKey,
Results: results,
Statuses: statuses,
Messages: shim.Messages,
}
return req, nil
+25 -11
View File
@@ -14,29 +14,43 @@ x-default-settings:
soft: 1000000000
services:
ubuntu14-osquery:
image: "kolide/osquery:${KOLIDE_OSQUERY_VERSION}"
volumes: *default-volumes
environment: *default-environment
command: *default-command
ulimits: *default-ulimits
ubuntu16-osquery:
image: "kolide/ubuntu16-osquery:${KOLIDE_OSQUERY_VERSION}"
image: "dactiv/osquery:4.3.0-ubuntu16.04"
volumes: *default-volumes
environment: *default-environment
command: *default-command
ulimits: *default-ulimits
centos7-osquery:
image: "kolide/centos7-osquery:${KOLIDE_OSQUERY_VERSION}"
ubuntu18-osquery:
image: "dactiv/osquery:4.3.0-ubuntu18.04"
volumes: *default-volumes
environment: *default-environment
command: *default-command
ulimits: *default-ulimits
ubuntu20-osquery:
image: "dactiv/osquery:4.3.0-ubuntu20.04"
volumes: *default-volumes
environment: *default-environment
command: *default-command
ulimits: *default-ulimits
centos6-osquery:
image: "kolide/centos6-osquery:${KOLIDE_OSQUERY_VERSION}"
image: "dactiv/osquery:4.3.0-centos6"
volumes: *default-volumes
environment: *default-environment
command: *default-command
ulimits: *default-ulimits
centos7-osquery:
image: "dactiv/osquery:4.3.0-centos7"
volumes: *default-volumes
environment: *default-environment
command: *default-command
ulimits: *default-ulimits
centos8-osquery:
image: "dactiv/osquery:4.3.0-centos8"
volumes: *default-volumes
environment: *default-environment
command: *default-command