| 22 | const INITIAL_QUEUE_SIZE = 5 |
| 23 | |
| 24 | func main() { |
| 25 | extensionName := path.Base(os.Args[0]) |
| 26 | printPrefix := fmt.Sprintf("[%s]", extensionName) |
| 27 | logger := log.WithFields(log.Fields{"agent": extensionName}) |
| 28 | |
| 29 | extensionClient := extension.NewClient(os.Getenv("AWS_LAMBDA_RUNTIME_API")) |
| 30 | |
| 31 | ctx, cancel := context.WithCancel(context.Background()) |
| 32 | |
| 33 | sigs := make(chan os.Signal, 1) |
| 34 | signal.Notify(sigs, syscall.SIGTERM, syscall.SIGINT) |
| 35 | go func() { |
| 36 | s := <-sigs |
| 37 | cancel() |
| 38 | logger.Info(printPrefix, "Received", s) |
| 39 | logger.Info(printPrefix, "Exiting") |
| 40 | }() |
| 41 | |
| 42 | // Register extension as soon as possible |
| 43 | _, err := extensionClient.Register(ctx, extensionName) |
| 44 | if err != nil { |
| 45 | panic(err) |
| 46 | } |
| 47 | |
| 48 | // Create S3 Logger |
| 49 | logsApiLogger, err := agent.NewS3Logger() |
| 50 | if err != nil { |
| 51 | logger.Fatal(err) |
| 52 | } |
| 53 | |
| 54 | // A synchronous queue that is used to put logs from the goroutine (producer) |
| 55 | // and process the logs from main goroutine (consumer) |
| 56 | logQueue := queue.New(INITIAL_QUEUE_SIZE) |
| 57 | // Helper function to empty the log queue |
| 58 | var logsStr string = "" |
| 59 | flushLogQueue := func(force bool) { |
| 60 | for !(logQueue.Empty() && (force || strings.Contains(logsStr, string(logsapi.RuntimeDone)))) { |
| 61 | logs, err := logQueue.Get(1) |
| 62 | if err != nil { |
| 63 | logger.Error(printPrefix, err) |
| 64 | return |
| 65 | } |
| 66 | logsStr = fmt.Sprintf("%v", logs[0]) |
| 67 | err = logsApiLogger.PushLog(logsStr) |
| 68 | if err != nil { |
| 69 | logger.Error(printPrefix, err) |
| 70 | return |
| 71 | } |
| 72 | } |
| 73 | } |
| 74 | |
| 75 | // Create Logs API agent |
| 76 | logsApiAgent, err := agent.NewHttpAgent(logsApiLogger, logQueue) |
| 77 | if err != nil { |
| 78 | logger.Fatal(err) |
| 79 | } |
| 80 | |
| 81 | // Subscribe to logs API |