Class: Fluent::GoogleCloudOutput

Inherits:
BufferedOutput
  • Object
show all
Defined in:
lib/fluent/plugin/out_google_cloud.rb

Constant Summary collapse

APPENGINE_SERVICE =

Constants for Google service names.

'appengine.googleapis.com'
COMPUTE_SERVICE =
'compute.googleapis.com'
DATAFLOW_SERVICE =
'dataflow.googleapis.com'

Instance Attribute Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize ⇒ GoogleCloudOutput

Returns a new instance of GoogleCloudOutput.



52
53
54
55
56
57
58
# File 'lib/fluent/plugin/out_google_cloud.rb', line 52

def initialize
  super
  require 'cgi'
  require 'google/api_client'
  require 'google/api_client/auth/compute_service_account'
  require 'open-uri'
end

Instance Attribute Details

#common_labels ⇒ Object (readonly)

Returns the value of attribute common_labels.



50
51
52
# File 'lib/fluent/plugin/out_google_cloud.rb', line 50

def common_labels
  @common_labels
end

#gae_backend_name ⇒ Object (readonly)

Returns the value of attribute gae_backend_name.



47
48
49
# File 'lib/fluent/plugin/out_google_cloud.rb', line 47

def gae_backend_name
  @gae_backend_name
end

#gae_backend_version ⇒ Object (readonly)

Returns the value of attribute gae_backend_version.



48
49
50
# File 'lib/fluent/plugin/out_google_cloud.rb', line 48

def gae_backend_version
  @gae_backend_version
end

#project_id ⇒ Object (readonly)

Expose attr_readers to make testing of metadata more direct than only testing it indirectly through metadata sent with logs.



43
44
45
# File 'lib/fluent/plugin/out_google_cloud.rb', line 43

def project_id
  @project_id
end

#running_on_managed_vm ⇒ Object (readonly)

Returns the value of attribute running_on_managed_vm.



46
47
48
# File 'lib/fluent/plugin/out_google_cloud.rb', line 46

def running_on_managed_vm
  @running_on_managed_vm
end

#service_name ⇒ Object (readonly)

Returns the value of attribute service_name.



49
50
51
# File 'lib/fluent/plugin/out_google_cloud.rb', line 49

def service_name
  @service_name
end

#vm_id ⇒ Object (readonly)

Returns the value of attribute vm_id.



45
46
47
# File 'lib/fluent/plugin/out_google_cloud.rb', line 45

def vm_id
  @vm_id
end

#zone ⇒ Object (readonly)

Returns the value of attribute zone.



44
45
46
# File 'lib/fluent/plugin/out_google_cloud.rb', line 44

def zone
  @zone
end

Instance Method Details

#configure(conf) ⇒ Object



60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
# File 'lib/fluent/plugin/out_google_cloud.rb', line 60

def configure(conf)
  super

  case @auth_method
  when 'private_key'
    if !@private_key_email
      raise Fluent::ConfigError, ('"private_key_email" must be ' +
          'specified if auth_method is "private_key"')
    elsif !@private_key_path
      raise Fluent::ConfigError, ('"private_key_path" must be ' +
          'specified if auth_method is "private_key"')
    elsif !@private_key_passphrase
      raise Fluent::ConfigError, ('"private_key_passphrase" must be ' +
          'specified if auth_method is "private_key"')
    end
  when 'compute_engine_service_account'
    # pass
  else
    raise Fluent::ConfigError,
        ('Unrecognized "auth_method" parameter. Please specify either ' +
         '"compute_engine_service_account" or "private_key".')
  end
end

#format(tag, time, record) ⇒ Object



134
135
136
# File 'lib/fluent/plugin/out_google_cloud.rb', line 134

def format(tag, time, record)
  [tag, time, record].to_msgpack
end

#shutdown ⇒ Object



130
131
132
# File 'lib/fluent/plugin/out_google_cloud.rb', line 130

def shutdown
  super
end

#start ⇒ Object



84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
# File 'lib/fluent/plugin/out_google_cloud.rb', line 84

def start
  super

  init_api_client()

  @successful_call = false

  # Grab metadata about the Google Compute Engine instance that we're on.
  @project_id = ('project/project-id')
  fully_qualified_zone = ('instance/zone')
  @zone = fully_qualified_zone.rpartition('/')[2]
  @vm_id = ('instance/id')
  # TODO: Send instance tags and/or hostname with the logs as well?
  @common_labels = {}

  # If this is running on a Managed VM, grab the relevant App Engine
  # metadata as well.
  # TODO: Add config options for these to allow for running outside GCE?
  attributes_string = ('instance/attributes/')
  attributes = attributes_string.split
  if (attributes.include?('gae_backend_name') &&
      attributes.include?('gae_backend_version'))
    @running_on_managed_vm = true
    @gae_backend_name =
        ('instance/attributes/gae_backend_name')
    @gae_backend_version =
        ('instance/attributes/gae_backend_version')
    @service_name = APPENGINE_SERVICE
    common_labels["#{APPENGINE_SERVICE}/module_id"] = @gae_backend_name
    common_labels["#{APPENGINE_SERVICE}/version_id"] = @gae_backend_version
  elsif (attributes.include?('job_id'))
    @running_on_managed_vm = false
    @service_name = DATAFLOW_SERVICE
    @dataflow_job_id = ('instance/attributes/job_id')
    common_labels["#{DATAFLOW_SERVICE}/job_id"] = @dataflow_job_id
  else
    @running_on_managed_vm = false
    @service_name = COMPUTE_SERVICE
  end

  if (@service_name != DATAFLOW_SERVICE)
    common_labels["#{COMPUTE_SERVICE}/resource_type"] = 'instance'
    common_labels["#{COMPUTE_SERVICE}/resource_id"] = @vm_id
  end
end

#write(chunk) ⇒ Object



138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
# File 'lib/fluent/plugin/out_google_cloud.rb', line 138

def write(chunk)
  # Group the entries since we have to make one call per tag.
  grouped_entries = {}
  chunk.msgpack_each do |tag, *arr|
    if !grouped_entries.has_key?(tag)
      grouped_entries[tag] = []
    end
    grouped_entries[tag].push(arr)
  end

  grouped_entries.each do |tag, arr|
    write_log_entries_request = {
      'commonLabels' => @common_labels,
      'entries' => [],
    }
    arr.each do |time, record|
      if (record.has_key?('timeNanos'))
        ts_secs = (record['timeNanos'] / 1000000000).to_i
        ts_nanos = record['timeNanos'] % 1000000000
        record.delete('timeNanos')
      else
        timestamp = Time.at(time)
        ts_secs = timestamp.tv_sec
        ts_nanos = timestamp.tv_nsec
      end
      entry = {
        'metadata' => {
          'serviceName' => @service_name,
          'projectId' => @project_id,
          'zone' => @zone,
          'timestamp' => {
            'seconds' => ts_secs,
            'nanos' => ts_nanos
          },
        },
      }
      if record.has_key?('severity')
        entry['metadata']['severity'] = parse_severity(record['severity'])
        record.delete('severity')
      else
        entry['metadata']['severity'] = 'DEFAULT'
      end

      # use textPayload if the only remainaing key is 'message',
      # otherwise use a struct.
      if (record.size == 1 && record.has_key?('message'))
        entry['textPayload'] = record['message']
      else
        entry['structPayload'] = record
      end
      write_log_entries_request['entries'].push(entry)
    end

    # Add a prefix to VMEngines logs to prevent namespace collisions,
    # and also escape the log name.
    log_name = CGI::escape(@running_on_managed_vm ?
                           "#{APPENGINE_SERVICE}/#{tag}" : tag)
    url = ('https://logging.googleapis.com/v1beta3/projects/' +
           "#{@project_id}/logs/#{log_name}/entries:write")
    begin
      client = api_client()
      request = client.generate_request({
        :uri => url,
        :body_object => write_log_entries_request,
        :http_method => 'POST',
        :authenticated => true
      })
      client.execute!(request)
      # Let the user explicitly know when the first call succeeded,
      # to aid with verification and troubleshooting.
      if (!@successful_call)
        @successful_call = true
        $log.info "Successfully sent to Google Cloud Logging API."
      end
    # Allow most exceptions to propagate, which will cause fluentd to
    # retry (with backoff), but in some cases we catch the error and
    # drop the request (we will emit a log message in those cases).
    rescue Google::APIClient::ClientError => error
      # Most ClientErrors indicate a problem with the request itself and
      # should not be retried, unless it is an authentication issue, in
      # which case we will retry the request via re-raising the exception.
      if (is_retriable_client_error(error))
        raise error
      end
      log_write_failure(write_log_entries_request, error)
    rescue JSON::GeneratorError => error
      # This happens if the request contains illegal characters;
      # do not retry it because it will fail repeatedly.
      log_write_failure(write_log_entries_request, error)
    end
  end
end