From 994071243fc80990cf09e820b671006e9bd31c76 Mon Sep 17 00:00:00 2001 From: Gab <24553253+gabrix73@users.noreply.github.com> Date: Thu, 13 Aug 2026 14:04:30 +0200 Subject: Separate active YAMN code from Katzenpost PoC --- katzenpost/cmd/yamn-dispatcher/main.go | 115 +++++++++++++++++++++++++++++++++ katzenpost/cmd/yamn-submit/main.go | 73 +++++++++++++++++++++ 2 files changed, 188 insertions(+) create mode 100644 katzenpost/cmd/yamn-dispatcher/main.go create mode 100644 katzenpost/cmd/yamn-submit/main.go (limited to 'katzenpost/cmd') diff --git a/katzenpost/cmd/yamn-dispatcher/main.go b/katzenpost/cmd/yamn-dispatcher/main.go new file mode 100644 index 0000000..8f10034 --- /dev/null +++ b/katzenpost/cmd/yamn-dispatcher/main.go @@ -0,0 +1,115 @@ +package main + +import ( + "errors" + "flag" + "fmt" + "os" + "path/filepath" + "strings" + "time" + + "git.virebent.art/virebent/yamnweb/katzenpost/dispatcher" + "github.com/katzenpost/katzenpost/core/log" + "github.com/katzenpost/katzenpost/server/cborplugin" +) + +const capability = "yamn-dispatch-v1" + +type plugin struct { + write func(cborplugin.Command) + reassembler *dispatcher.Reassembler + allowedDomains map[string]struct{} +} + +func (p *plugin) OnCommand(command cborplugin.Command) error { + request, ok := command.(*cborplugin.Request) + if !ok { + return errors.New("unexpected plugin command") + } + frame, err := dispatcher.DecodeFrame(request.Payload) + if err != nil { + return err + } + encoded, complete, err := p.reassembler.Add(frame) + if err != nil || !complete { + return err + } + envelope, err := dispatcher.DecodeEnvelope(encoded) + if err != nil { + return err + } + if err := envelope.Validate(p.allowedDomains); err != nil { + return err + } + // PoC sink: successful validation intentionally has no external side effect. + return nil +} + +func (p *plugin) RegisterConsumer(server *cborplugin.Server) { + p.write = server.Write +} + +func parseDomains(value string) (map[string]struct{}, error) { + domains := make(map[string]struct{}) + for _, domain := range strings.Split(value, ",") { + domain = strings.ToLower(strings.TrimSpace(domain)) + if domain == "" || strings.ContainsAny(domain, "@/\\: ") { + return nil, errors.New("invalid allowed domain") + } + domains[domain] = struct{}{} + } + if len(domains) == 0 { + return nil, errors.New("at least one allowed domain is required") + } + return domains, nil +} + +func run() error { + var allowed string + var logDir string + var logLevel string + flag.StringVar(&allowed, "allowed-domains", "remailer.example", "comma-separated entry remailer domains") + flag.StringVar(&logDir, "log-dir", "/tmp", "operational log directory") + flag.StringVar(&logLevel, "log-level", "NOTICE", "Katzenpost log level") + flag.Parse() + + domains, err := parseDomains(allowed) + if err != nil { + return err + } + info, err := os.Stat(logDir) + if err != nil || !info.IsDir() { + return errors.New("log directory is unavailable") + } + backend, err := log.New(filepath.Join(logDir, "yamn-dispatch.log"), logLevel, false) + if err != nil { + return fmt.Errorf("initialize logging: %w", err) + } + logger := backend.GetLogger("yamn_dispatch") + + socketDir, err := os.MkdirTemp("", "yamn-dispatch-") + if err != nil { + return fmt.Errorf("create socket directory: %w", err) + } + defer os.RemoveAll(socketDir) + socketPath := filepath.Join(socketDir, "plugin.sock") + service := &plugin{ + reassembler: dispatcher.NewReassembler(5 * time.Minute), + allowedDomains: domains, + } + server := cborplugin.NewServer(logger, socketPath, new(cborplugin.RequestFactory), service) + if _, err := fmt.Fprintln(os.Stdout, socketPath); err != nil { + return fmt.Errorf("publish socket path: %w", err) + } + server.Accept() + server.Wait() + return nil +} + +func main() { + if err := run(); err != nil { + fmt.Fprintln(os.Stderr, "yamn-dispatcher failed") + os.Exit(1) + } +} diff --git a/katzenpost/cmd/yamn-submit/main.go b/katzenpost/cmd/yamn-submit/main.go new file mode 100644 index 0000000..cde35f9 --- /dev/null +++ b/katzenpost/cmd/yamn-submit/main.go @@ -0,0 +1,73 @@ +package main + +import ( + "encoding/json" + "errors" + "flag" + "fmt" + "io" + "os" + "time" + + "git.virebent.art/virebent/yamnweb/katzenpost/dispatcher" + "github.com/katzenpost/hpqc/hash" + clientconfig "github.com/katzenpost/katzenpost/client/config" + "github.com/katzenpost/katzenpost/client/thin" +) + +const capability = "yamn-dispatch-v1" + +func run() error { + var configPath string + var settle time.Duration + flag.StringVar(&configPath, "config", "thinclient.toml", "thin-client configuration") + flag.DurationVar(&settle, "settle", 2*time.Second, "time allowed for kpclientd to queue frames") + flag.Parse() + + input, err := io.ReadAll(io.LimitReader(os.Stdin, dispatcher.MaxEnvelopeBytes*2)) + if err != nil { + return fmt.Errorf("read envelope: %w", err) + } + var envelope dispatcher.Envelope + if err := json.Unmarshal(input, &envelope); err != nil { + return errors.New("invalid envelope input") + } + encoded, err := dispatcher.EncodeEnvelope(envelope) + if err != nil { + return err + } + frames, err := dispatcher.SplitEnvelope(encoded) + if err != nil { + return err + } + + config, err := thin.LoadFile(configPath) + if err != nil { + return fmt.Errorf("load thin-client config: %w", err) + } + client := thin.NewThinClient(config, &clientconfig.Logging{Level: "ERROR", Disable: true}) + defer client.Close() + if err := client.Dial(); err != nil { + return fmt.Errorf("connect to kpclientd: %w", err) + } + service, err := client.GetService(capability) + if err != nil { + return fmt.Errorf("find dispatcher service: %w", err) + } + destination := hash.Sum256(service.MixDescriptor.IdentityKey) + for _, frame := range frames { + if err := client.SendMessageWithoutReply(frame, &destination, service.RecipientQueueID); err != nil { + return fmt.Errorf("send frame: %w", err) + } + } + time.Sleep(settle) + fmt.Fprintln(os.Stdout, "accepted") + return nil +} + +func main() { + if err := run(); err != nil { + fmt.Fprintf(os.Stderr, "submission failed: %v\n", err) + os.Exit(1) + } +} -- cgit v1.2.3