diff --git a/server/kolide/osquery.go b/server/kolide/osquery.go index 14c0382811..cc0e2d5630 100644 --- a/server/kolide/osquery.go +++ b/server/kolide/osquery.go @@ -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) diff --git a/server/launcher/launcher.go b/server/launcher/launcher.go index 7dd46b7aeb..d71f6abecf 100644 --- a/server/launcher/launcher.go +++ b/server/launcher/launcher.go @@ -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") } diff --git a/server/launcher/launcher_test.go b/server/launcher/launcher_test.go index 07a4ec40b8..8cdbb3eae7 100644 --- a/server/launcher/launcher_test.go +++ b/server/launcher/launcher_test.go @@ -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 }, diff --git a/server/mock/service_osquery.go b/server/mock/service_osquery.go index 54e2f63f9c..8a20772f58 100644 --- a/server/mock/service_osquery.go +++ b/server/mock/service_osquery.go @@ -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) { diff --git a/server/service/endpoint_osquery.go b/server/service/endpoint_osquery.go index 4640d81fa4..f7bceca11e 100644 --- a/server/service/endpoint_osquery.go +++ b/server/service/endpoint_osquery.go @@ -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 } diff --git a/server/service/logging_osquery.go b/server/service/logging_osquery.go index 19dc0587f7..6a99668ec0 100644 --- a/server/service/logging_osquery.go +++ b/server/service/logging_osquery.go @@ -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 } diff --git a/server/service/service_osquery.go b/server/service/service_osquery.go index 54a1362a1a..141a35ec84 100644 --- a/server/service/service_osquery.go +++ b/server/service/service_osquery.go @@ -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} } diff --git a/server/service/service_osquery_test.go b/server/service/service_osquery_test.go index a0f9426a04..18cb6e9fcc 100644 --- a/server/service/service_osquery_test.go +++ b/server/service/service_osquery_test.go @@ -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) } diff --git a/server/service/transport_osquery.go b/server/service/transport_osquery.go index c7c28c930e..6208f0daba 100644 --- a/server/service/transport_osquery.go +++ b/server/service/transport_osquery.go @@ -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 diff --git a/tools/osquery/docker-compose.yml b/tools/osquery/docker-compose.yml index 61c84fd445..a69282bfbd 100644 --- a/tools/osquery/docker-compose.yml +++ b/tools/osquery/docker-compose.yml @@ -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