mirror of
https://github.com/psviderski/uncloud.git
synced 2026-08-26 19:13:34 +00:00
* Add server side of journal logs This add the server side and grpc methods to get a journal logs from a machine. It repurposes ServiceLogEntry for these logs to keep the changes somewhat to a minimum. And it lets us re-use the merging of the various logs. In the protobufs ContainerLog has been renamed to just Log and LogEntry, as these are now also used for journal logs. It does api.LogOptions in more places to reduce the various logOpts that were used. It does not yet plumb it through to the uc client, that needs a follow up pr. Following logs is also not yet implemented. Signed-off-by: Miek Gieben <miek@miek.nl> * Fix test too Signed-off-by: Miek Gieben <miek@miek.nl> * remove entire comment Signed-off-by: Miek Gieben <miek@miek.nl> * Implement the follow option, untested mind you Signed-off-by: Miek Gieben <miek@miek.nl> * update debug line Signed-off-by: Miek Gieben <miek@miek.nl> * internal/jounal: First batch of PR comments Signed-off-by: Miek Gieben <miek@miek.nl> * internal/journal: code review comments Signed-off-by: Miek Gieben <miek@miek.nl> * Manually apply suggestion Signed-off-by: Miek Gieben <miek@miek.nl> * apply comment manually Signed-off-by: Miek Gieben <miek@miek.nl> * internal/journal: add unit test Signed-off-by: Miek Gieben <miek@miek.nl> * Use testify Signed-off-by: Miek Gieben <miek@miek.nl> * -amFix scanner.Err checking Signed-off-by: Miek Gieben <miek@miek.nl> * Implement code review comments Signed-off-by: Miek Gieben <miek@miek.nl> --------- Signed-off-by: Miek Gieben <miek@miek.nl>
160 lines
4.4 KiB
Go
160 lines
4.4 KiB
Go
package client
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
|
|
"github.com/docker/docker/pkg/stringid"
|
|
"github.com/psviderski/uncloud/internal/machine/api/pb"
|
|
"github.com/psviderski/uncloud/pkg/api"
|
|
)
|
|
|
|
// ServiceLogs streams log entries from all service containers in chronological order based on timestamps.
|
|
// Keep in mind that perfect ordering of log events across multiple machines can't be guaranteed due to the
|
|
// imperfection of physical clocks or potential clock skew between machines.
|
|
// It uses a low watermark algorithm to ensure proper ordering across multiple machines.
|
|
// Heartbeat entries from the server advance the watermark to enable timely emission of buffered logs.
|
|
func (cli *Client) ServiceLogs(
|
|
ctx context.Context, serviceNameOrID string, opts api.ServiceLogsOptions,
|
|
) (api.Service, <-chan api.ServiceLogEntry, error) {
|
|
svc, err := cli.InspectService(ctx, serviceNameOrID)
|
|
if err != nil {
|
|
return svc, nil, fmt.Errorf("inspect service: %w", err)
|
|
}
|
|
|
|
if len(svc.Containers) == 0 {
|
|
return svc, nil, fmt.Errorf("no containers found for service: %s", serviceNameOrID)
|
|
}
|
|
|
|
machines, err := cli.ListMachines(ctx, &api.MachineFilter{
|
|
NamesOrIDs: opts.Machines,
|
|
})
|
|
if err != nil {
|
|
return svc, nil, fmt.Errorf("list machines: %w", err)
|
|
}
|
|
|
|
ctrStreams := make([]<-chan api.ServiceLogEntry, 0, len(svc.Containers))
|
|
for _, ctr := range svc.Containers {
|
|
// Skip containers not running on the specified machines.
|
|
m := machines.FindByNameOrID(ctr.MachineID)
|
|
if len(opts.Machines) > 0 && m == nil {
|
|
continue
|
|
}
|
|
|
|
// Machine name for ServiceLogEntry metadata and friendlier error message.
|
|
machineName := ctr.MachineID
|
|
if m != nil {
|
|
machineName = m.Machine.Name
|
|
}
|
|
|
|
stream, err := cli.ContainerLogs(ctx, ctr.MachineID, ctr.Container.ID, opts)
|
|
if err != nil {
|
|
return svc, nil, fmt.Errorf("stream logs from service container '%s' on machine '%s': %w",
|
|
stringid.TruncateID(ctr.Container.ID), machineName, err)
|
|
}
|
|
|
|
// Enrich log entries from the container with service metadata.
|
|
metadata := api.ServiceLogEntryMetadata{
|
|
ServiceID: svc.ID,
|
|
ServiceName: svc.Name,
|
|
ContainerID: ctr.Container.ID,
|
|
MachineID: ctr.MachineID,
|
|
MachineName: machineName,
|
|
}
|
|
enrichedStream := logsStreamWithServiceMetadata(stream, metadata)
|
|
ctrStreams = append(ctrStreams, enrichedStream)
|
|
}
|
|
|
|
if len(ctrStreams) == 0 {
|
|
return svc, nil, errors.New("no service containers found on the specified machine(s)")
|
|
}
|
|
|
|
// Use the log merger to combine streams from all containers in chronological order.
|
|
merger := NewLogMerger(ctrStreams, DefaultLogMergerOptions)
|
|
mergedStream := merger.Stream()
|
|
|
|
return svc, mergedStream, nil
|
|
}
|
|
|
|
// ContainerLogs streams log entries from a single container on a specified machine.
|
|
func (cli *Client) ContainerLogs(
|
|
ctx context.Context, machineNameOrID string, containerID string, opts api.ServiceLogsOptions,
|
|
) (<-chan api.LogEntry, error) {
|
|
proxyCtx, _, err := cli.ProxyMachinesContext(ctx, []string{machineNameOrID})
|
|
if err != nil {
|
|
return nil, fmt.Errorf("create request context to proxy to machine '%s': %w", machineNameOrID, err)
|
|
}
|
|
|
|
req := &pb.LogsRequest{
|
|
Id: containerID,
|
|
Follow: opts.Follow,
|
|
Tail: int32(opts.Tail),
|
|
Since: opts.Since,
|
|
Until: opts.Until,
|
|
}
|
|
if !opts.Follow && opts.Tail == 0 {
|
|
// If not following and tail is 0, set tail to -1 to return all logs.
|
|
// Otherwise, no logs will be returned at all.
|
|
req.Tail = -1
|
|
}
|
|
|
|
stream, err := cli.Docker.GRPCClient.ContainerLogs(proxyCtx, req)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
ch := make(chan api.LogEntry)
|
|
|
|
go func() {
|
|
defer close(ch)
|
|
|
|
for {
|
|
pbEntry, err := stream.Recv()
|
|
if err == io.EOF {
|
|
return
|
|
}
|
|
if err != nil {
|
|
ch <- api.LogEntry{
|
|
Err: err,
|
|
}
|
|
return
|
|
}
|
|
|
|
entry := api.LogEntry{
|
|
Stream: api.LogStreamTypeFromProto(pbEntry.Stream),
|
|
Message: pbEntry.Message,
|
|
Timestamp: pbEntry.Timestamp.AsTime(),
|
|
}
|
|
|
|
select {
|
|
case ch <- entry:
|
|
case <-ctx.Done():
|
|
return
|
|
}
|
|
}
|
|
}()
|
|
|
|
return ch, nil
|
|
}
|
|
|
|
// logsStreamWithServiceMetadata wraps a container logs stream and enriches each log entry with service metadata.
|
|
func logsStreamWithServiceMetadata(
|
|
stream <-chan api.LogEntry, metadata api.ServiceLogEntryMetadata,
|
|
) <-chan api.ServiceLogEntry {
|
|
out := make(chan api.ServiceLogEntry)
|
|
|
|
go func() {
|
|
for entry := range stream {
|
|
out <- api.ServiceLogEntry{
|
|
Metadata: metadata,
|
|
LogEntry: entry,
|
|
}
|
|
}
|
|
close(out)
|
|
}()
|
|
|
|
return out
|
|
}
|