file.go 2.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116
  1. package dockerlogcollector
  2. import (
  3. "context"
  4. "time"
  5. "strings"
  6. "path"
  7. "io/ioutil"
  8. "go-common/library/log"
  9. "go-common/app/service/ops/log-agent/pipeline"
  10. "github.com/docker/docker/api/types"
  11. "github.com/docker/docker/client"
  12. )
  13. type DockerLogCollector struct {
  14. c *Config
  15. client *client.Client
  16. ctx context.Context
  17. cancel context.CancelFunc
  18. }
  19. type configItem struct {
  20. configPath string
  21. MergedDir string
  22. }
  23. func InitDockerLogCollector(ctx context.Context, c *Config) (err error) {
  24. if err = c.ConfigValidate(); err != nil {
  25. return err
  26. }
  27. collector := new(DockerLogCollector)
  28. collector.c = c
  29. collector.ctx, collector.cancel = context.WithCancel(ctx)
  30. // init docker client
  31. collector.client, err = client.NewEnvClient()
  32. if err != nil {
  33. return err
  34. }
  35. go collector.scan()
  36. return nil
  37. }
  38. func (collector *DockerLogCollector) getConfigs() ([]*configItem, error) {
  39. var (
  40. configItems = make([]*configItem, 0)
  41. mergedDir string
  42. ok bool
  43. )
  44. containers, err := collector.client.ContainerList(collector.ctx, types.ContainerListOptions{})
  45. if err != nil {
  46. return nil, err
  47. }
  48. for _, container := range containers {
  49. info, err := collector.client.ContainerInspect(collector.ctx, container.ID)
  50. if err != nil {
  51. log.Error("failed to inspect container: %s", container.ID)
  52. continue
  53. }
  54. // get overlay2 info
  55. if info.GraphDriver.Name != "overlay2" {
  56. log.Error("only overlay2 is supported")
  57. continue
  58. }
  59. mergedDir, ok = info.GraphDriver.Data["MergedDir"]
  60. if !ok {
  61. log.Error("failed to get MergedDir of container:%s", container.ID)
  62. }
  63. for _, env := range info.Config.Env {
  64. if strings.HasPrefix(env, collector.c.
  65. ConfigEnv) {
  66. for _, path := range strings.Split(strings.TrimPrefix(env, collector.c.ConfigEnv+"="), ",") {
  67. configItems = append(configItems, &configItem{path, mergedDir})
  68. }
  69. }
  70. }
  71. }
  72. return configItems, nil
  73. }
  74. func (collector *DockerLogCollector) scan() {
  75. ticker := time.Tick(time.Duration(collector.c.ScanInterval))
  76. for {
  77. select {
  78. case <-ticker:
  79. configItems, err := collector.getConfigs()
  80. if err != nil {
  81. log.Error("failed to scan hostlogcollector config file list: %s", err)
  82. continue
  83. }
  84. for _, item := range configItems {
  85. configPath := path.Join(item.MergedDir, item.configPath)
  86. config, err := ioutil.ReadFile(configPath)
  87. if err != nil {
  88. log.Error("filed to read hostlogcollector config file %s: %s", configPath, err)
  89. continue
  90. }
  91. if !pipeline.PipelineManagement.PipelineExisted(configPath) {
  92. ctx := context.WithValue(collector.ctx, "MergedDir", item.MergedDir)
  93. go pipeline.PipelineManagement.StartPipeline(ctx, configPath, string(config))
  94. }
  95. }
  96. case <-collector.ctx.Done():
  97. return
  98. }
  99. }
  100. }