gcpdiag.queries.logs

Queries related to Cloud Logging. The main functionality is querying log entries, which is supposed to be used as follows: 1. Call query() with the logs query parameters that you need. This returns a LogsQuery object which can be used to retrieve the logs later. 2. Call execute_queries() to execute all log query jobs. Similar queries will be grouped together to minimize the number of required API calls. Multiple queries will be done in parallel, while always respecting the Cloud Logging limit of 60 queries per 60 seconds. 3. Use the entries property on the LogsQuery object to iterate over the fetched logs. Note that the entries are not guaranteed to be filtered by what was given in the "filter_str" argument to query(), you will need to filter out the entries in code as well when iterating over the log entries. Side note: this module is not called 'logging' to avoid using the same name as the standard python library for logging.
class LogsQuery:
67class LogsQuery:
68  """A log search job that was started with prefetch_logs()."""
69
70  job: _LogsQueryJob
71
72  def __init__(self, job):
73    self.job = job
74
75  @property
76  def entries(self) -> Sequence:
77    if not self.job.future:
78      raise RuntimeError("log query wasn't executed. did you forget to call execute_queries()?")
79    elif self.job.future.running():
80      logging.debug(
81        'waiting for logs query results (project: %s, resource type: %s)',
82        self.job.project_id,
83        self.job.resource_type,
84      )
85    return self.job.future.result()

A log search job that was started with prefetch_logs().

LogsQuery(job)
72  def __init__(self, job):
73    self.job = job
job: gcpdiag.queries.logs._LogsQueryJob
entries: Sequence
75  @property
76  def entries(self) -> Sequence:
77    if not self.job.future:
78      raise RuntimeError("log query wasn't executed. did you forget to call execute_queries()?")
79    elif self.job.future.running():
80      logging.debug(
81        'waiting for logs query results (project: %s, resource type: %s)',
82        self.job.project_id,
83        self.job.resource_type,
84      )
85    return self.job.future.result()
jobs_todo: Dict[Tuple[str, str, str], gcpdiag.queries.logs._LogsQueryJob] = {}
class LogEntryShort:
 91class LogEntryShort:
 92  """A common log entry"""
 93
 94  _text: str
 95  _timestamp: Optional[datetime.datetime]
 96
 97  def __init__(self, raw_entry):
 98    if isinstance(raw_entry, dict):
 99      self._text = get_path(raw_entry, ('textPayload',), default='')
100      self._timestamp = log_entry_timestamp(raw_entry)
101
102    if isinstance(raw_entry, str):
103      self._text = raw_entry
104      # we could extract timestamp from serial entries
105      # but they are not always present
106      # and may be unreliable as we don't know the system clock setting
107      self._timestamp = None
108
109  @property
110  def text(self):
111    return self._text
112
113  @property
114  def timestamp(self):
115    return self._timestamp
116
117  @property
118  def timestamp_iso(self):
119    if self._timestamp:
120      return self._timestamp.astimezone().isoformat(sep=' ', timespec='seconds')
121    return None

A common log entry

LogEntryShort(raw_entry)
 97  def __init__(self, raw_entry):
 98    if isinstance(raw_entry, dict):
 99      self._text = get_path(raw_entry, ('textPayload',), default='')
100      self._timestamp = log_entry_timestamp(raw_entry)
101
102    if isinstance(raw_entry, str):
103      self._text = raw_entry
104      # we could extract timestamp from serial entries
105      # but they are not always present
106      # and may be unreliable as we don't know the system clock setting
107      self._timestamp = None
text
109  @property
110  def text(self):
111    return self._text
timestamp
113  @property
114  def timestamp(self):
115    return self._timestamp
timestamp_iso
117  @property
118  def timestamp_iso(self):
119    if self._timestamp:
120      return self._timestamp.astimezone().isoformat(sep=' ', timespec='seconds')
121    return None
class LogExclusion(gcpdiag.models.Resource):
124class LogExclusion(models.Resource):
125  """A log exclusion entry"""
126
127  _resource_data: dict
128  project_id: str
129
130  def __init__(self, project_id: str, resource_data: dict):
131    super().__init__(project_id)
132    self._resource_data = resource_data
133
134  @property
135  def full_path(self) -> str:
136    return self._resource_data['name']
137
138  @property
139  def filter(self) -> str:
140    return self._resource_data['filter']
141
142  @property
143  def disabled(self) -> bool:
144    if 'disabled' in self._resource_data:
145      return self._resource_data['disabled']
146    return False

A log exclusion entry

LogExclusion(project_id: str, resource_data: dict)
130  def __init__(self, project_id: str, resource_data: dict):
131    super().__init__(project_id)
132    self._resource_data = resource_data
project_id: str
264  @property
265  def project_id(self) -> str:
266    """Project id (not project number)."""
267    return self._project_id

Project id (not project number).

full_path: str
134  @property
135  def full_path(self) -> str:
136    return self._resource_data['name']

Returns the full path of this resource.

Example: 'projects/gcpdiag-gke-1-9b90/zones/europe-west4-a/clusters/gke1'

filter: str
138  @property
139  def filter(self) -> str:
140    return self._resource_data['filter']
disabled: bool
142  @property
143  def disabled(self) -> bool:
144    if 'disabled' in self._resource_data:
145      return self._resource_data['disabled']
146    return False
def query( project_id: str, resource_type: str, log_name: str, filter_str: str) -> LogsQuery:
149def query(project_id: str, resource_type: str, log_name: str, filter_str: str) -> LogsQuery:
150  # Aggregate by project_id, resource_type, log_name
151  job_key = (project_id, resource_type, log_name)
152  job = jobs_todo.setdefault(
153    job_key,
154    _LogsQueryJob(
155      project_id=project_id,
156      resource_type=resource_type,
157      log_name=log_name,
158      filters=set(),
159    ),
160  )
161  job.filters.add(filter_str)
162  return LogsQuery(job=job)
@caching.cached_api_call
def realtime_query(project_id, filter_str, start_time, end_time, disable_paging=False):
262@caching.cached_api_call
263def realtime_query(project_id, filter_str, start_time, end_time, disable_paging=False):
264  """Intended for use in only runbooks. use logs.query() for lint rules."""
265  logging_api = apis.get_api('logging', 'v2', project_id)
266
267  filter_lines = [filter_str]
268  filter_lines.append('timestamp>"%s"' % start_time.isoformat(timespec='seconds'))
269  filter_lines.append('timestamp<"%s"' % end_time.isoformat(timespec='seconds'))
270  filter_str = '\n'.join(filter_lines)
271  logging.debug(
272    'searching logs in project %s for logs between %s and %s',
273    project_id,
274    str(start_time),
275    str(end_time),
276  )
277  deque = Deque()
278  req = logging_api.entries().list(
279    body={
280      'resourceNames': [f'projects/{project_id}'],
281      'filter': filter_str,
282      'orderBy': 'timestamp desc',
283      'pageSize': config.get('logging_page_size'),
284    }
285  )
286  fetched_entries_count = 0
287  query_pages = 0
288  query_start_time = datetime.datetime.now()
289  while req is not None:
290    query_pages += 1
291    res = _ratelimited_execute(req)
292    if 'entries' in res:
293      for e in res['entries']:
294        fetched_entries_count += 1
295        deque.appendleft(e)
296
297    # Verify that we aren't above limits, exit otherwise.
298    if fetched_entries_count > config.get('logging_fetch_max_entries'):
299      logging.warning(
300        'maximum number of log entries (%d) reached (project: %s, query: %s).',
301        config.get('logging_fetch_max_entries'),
302        project_id,
303        filter_str.replace('\n', ' AND '),
304      )
305      return deque
306    run_time = (datetime.datetime.now() - query_start_time).total_seconds()
307    if run_time >= config.get('logging_fetch_max_time_seconds'):
308      logging.warning(
309        'maximum query runtime for log query reached (project: %s, query: %s).',
310        project_id,
311        filter_str.replace('\n', ' AND '),
312      )
313      return deque
314    if disable_paging:
315      break
316    req = logging_api.entries().list_next(req, res)
317    if req is not None:
318      logging.debug(
319        'still fetching logs (project: %s, max wait: %ds)',
320        project_id,
321        config.get('logging_fetch_max_time_seconds') - run_time,
322      )
323
324  query_end_time = datetime.datetime.now()
325  logging.debug(
326    'logging query run time: %s, pages: %d, query: %s',
327    query_end_time - query_start_time,
328    query_pages,
329    filter_str.replace('\n', ' AND '),
330  )
331
332  return deque

Intended for use in only runbooks. use logs.query() for lint rules.

def execute_queries( query_executor: gcpdiag.executor.ContextAwareExecutor, context: gcpdiag.models.Context):
335def execute_queries(query_executor: executor.ContextAwareExecutor, context: models.Context):
336  global jobs_todo
337  jobs_executing = jobs_todo
338  jobs_todo = {}
339  for job in jobs_executing.values():
340    job.future = query_executor.submit(_execute_query_job, job, context)
def log_entry_timestamp(log_entry: Mapping[str, Any]) -> datetime.datetime:
343def log_entry_timestamp(log_entry: Mapping[str, Any]) -> datetime.datetime:
344  # Use receiveTimestamp so that we don't have any time synchronization issues
345  # (i.e. don't trust the timestamp field)
346  timestamp = log_entry.get('receiveTimestamp', None)
347  if timestamp:
348    return dateutil.parser.parse(timestamp)
349  return timestamp
def format_log_entry(log_entry: dict) -> str:
352def format_log_entry(log_entry: dict) -> str:
353  """Format a log_entry, as returned by LogsQuery.entries to a simple one-line
354  string with the date and message."""
355  log_message = None
356  if 'jsonPayload' in log_entry:
357    for key in ['message', 'MESSAGE']:
358      if key in log_entry['jsonPayload']:
359        log_message = log_entry['jsonPayload'][key]
360        break
361  if log_message is None:
362    log_message = log_entry.get('textPayload')
363  log_date = log_entry_timestamp(log_entry)
364  log_date_str = log_date.astimezone().isoformat(sep=' ', timespec='seconds')
365  return f'{log_date_str}: {log_message}'

Format a log_entry, as returned by LogsQuery.entries to a simple one-line string with the date and message.

def exclusions(project_id: str) -> Optional[List[LogExclusion]]:
368def exclusions(project_id: str) -> Union[List[LogExclusion], None]:
369  logging_api = apis.get_api('logging', 'v2', project_id)
370  if not apis.is_enabled(project_id, 'logging'):
371    return None
372
373  log_exclusions: List[LogExclusion] = []
374
375  fetched_entries_count = 0
376  req = logging_api.exclusions().list(parent=f'projects/{project_id}')
377  while req is not None:
378    res = req.execute(num_retries=config.API_RETRIES)
379    fetched_entries_count += 1
380    if res:
381      for log_exclusion_resp in res['exclusions']:
382        log_exclusions.append(LogExclusion(project_id, log_exclusion_resp))
383    req = logging_api.exclusions().list_next(req, res)
384    if req is not None:
385      logging.debug(f'still fetching log exclusions for project {project_id}')
386  return log_exclusions