diff --git a/vmnet/events_darwin.go b/vmnet/events_darwin.go new file mode 100644 index 0000000..e371bcf --- /dev/null +++ b/vmnet/events_darwin.go @@ -0,0 +1,108 @@ +package vmnet + +/* +#include "vmnet_darwin.h" +*/ +import "C" +import ( + "fmt" + "runtime" + "runtime/cgo" + "sync" + "unsafe" +) + +// PacketsAvailableEventCallback receives the estimated number of packets ready to read. +type PacketsAvailableEventCallback func(estimatedCount int) + +type interfaceCallbackState struct { + mu sync.Mutex + iface unsafe.Pointer + queue unsafe.Pointer + handle cgo.Handle + stopped bool +} + +//export callPacketsAvailableEventCallback +func callPacketsAvailableEventCallback(handle C.uintptr_t, estimatedCount C.int) { + cgo.Handle(handle).Value().(PacketsAvailableEventCallback)(int(estimatedCount)) +} + +//export releasePacketsAvailableEventCallback +func releasePacketsAvailableEventCallback(handle C.uintptr_t) { + cgo.Handle(handle).Delete() +} + +// SetPacketsAvailableEventCallback registers one callback for packet availability. +// Pass nil to remove it. A callback may call this method to remove itself. +// Call Stop or remove the callback when it captures the Interface, so it can be collected. +func (i *Interface) SetPacketsAvailableEventCallback(callback PacketsAvailableEventCallback) error { + if i == nil || i.callbackState == nil { + return fmt.Errorf("interface is nil") + } + state := i.callbackState + state.mu.Lock() + defer state.mu.Unlock() + defer runtime.KeepAlive(i) + if state.stopped { + return fmt.Errorf("interface is stopped") + } + if callback == nil { + return state.clearLocked() + } + if state.queue != nil { + return fmt.Errorf("packets available callback is already set") + } + handle := cgo.NewHandle(callback) + var status C.uint32_t + queue := C.VmnetSetPacketsAvailableEventCallback(state.iface, C.uintptr_t(handle), &status) + if result := Return(status); result != ErrSuccess { + handle.Delete() + return fmt.Errorf("set packets available callback: %w", result) + } + state.queue = queue + state.handle = handle + return nil +} + +func (state *interfaceCallbackState) clearLocked() error { + if state.queue == nil { + return nil + } + result := Return(C.VmnetClearPacketsAvailableEventCallback(state.iface, state.queue, C.uintptr_t(state.handle))) + if result != ErrSuccess { + return fmt.Errorf("clear packets available callback: %w", result) + } + state.queue = nil + state.handle = 0 + return nil +} + +func (state *interfaceCallbackState) stop() error { + state.mu.Lock() + if state.stopped { + state.mu.Unlock() + return fmt.Errorf("interface is stopped") + } + if err := state.clearLocked(); err != nil { + state.mu.Unlock() + return err + } + state.stopped = true + state.mu.Unlock() + + result := Return(C.VmnetStopInterface(state.iface)) + if result != ErrSuccess { + state.mu.Lock() + state.stopped = false + state.mu.Unlock() + return fmt.Errorf("stop vmnet interface: %w", result) + } + return nil +} + +func (state *interfaceCallbackState) cleanup() { + state.mu.Lock() + defer state.mu.Unlock() + _ = state.clearLocked() +} diff --git a/vmnet/interface_darwin_test.go b/vmnet/interface_darwin_test.go index 0079962..2c98d8c 100644 --- a/vmnet/interface_darwin_test.go +++ b/vmnet/interface_darwin_test.go @@ -2,6 +2,7 @@ package vmnet_test import ( "testing" + "time" "github.com/Code-Hex/vz/v3/internal/osversion" "github.com/Code-Hex/vz/v3/vmnet" @@ -81,3 +82,74 @@ func TestStartInterfaceWithNilNetwork(t *testing.T) { t.Fatal("expected an error for a nil network") } } + +func TestPacketsAvailableEventCallback(t *testing.T) { + if err := osversion.MacOSAvailable(26); err != nil { + t.Skipf("vmnet interface requires macOS 26: %v", err) + } + config, err := vmnet.NewNetworkConfiguration(vmnet.HostMode) + if err != nil { + t.Fatal(err) + } + network, err := vmnet.NewNetwork(config) + if err != nil { + t.Fatal(err) + } + receiver, err := vmnet.StartInterfaceWithNetwork(network, nil) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = receiver.Stop() }) + sender, err := vmnet.StartInterfaceWithNetwork(network, nil) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = sender.Stop() }) + type eventResult struct { + count int + err error + } + available := make(chan eventResult, 1) + if err := receiver.SetPacketsAvailableEventCallback(func(count int) { + err := receiver.SetPacketsAvailableEventCallback(nil) + select { + case available <- eventResult{count, err}: + default: + } + }); err != nil { + t.Fatal(err) + } + if err := receiver.SetPacketsAvailableEventCallback(func(int) {}); err == nil { + t.Fatal("second callback registration must fail") + } + manager, err := vmnet.NewPktDescsManager(1, sender.MaxPacketSize) + if err != nil { + t.Fatal(err) + } + packet := make([]byte, 60) + for index := range 6 { + packet[index] = 0xff + } + packet[6] = 0x02 + packet[12], packet[13] = 0x08, 0x06 + if err := manager.SetPacket(0, packet); err != nil { + t.Fatal(err) + } + if err := sender.WritePackets(manager, 1); err != nil { + t.Fatal(err) + } + select { + case event := <-available: + if event.err != nil { + t.Fatal(event.err) + } + if event.count < 1 { + t.Fatalf("estimated packet count = %d", event.count) + } + case <-time.After(5 * time.Second): + t.Fatal("no packets available event") + } + if err := receiver.SetPacketsAvailableEventCallback(nil); err != nil { + t.Fatal(err) + } +} diff --git a/vmnet/vmnet_darwin.go b/vmnet/vmnet_darwin.go index b92024b..451a26d 100644 --- a/vmnet/vmnet_darwin.go +++ b/vmnet/vmnet_darwin.go @@ -483,6 +483,7 @@ type Interface struct { MaxPacketSize uint64 MaxReadPacketCount int MaxWritePacketCount int + callbackState *interfaceCallbackState } // StartInterfaceWithNetwork starts an Interface on a Network. @@ -514,25 +515,27 @@ func StartInterfaceWithNetwork(network *Network, interfaceDesc *xpc.Dictionary) MaxPacketSize: uint64(result.maxPacketSize), MaxReadPacketCount: int(result.maxReadPacketCount), MaxWritePacketCount: int(result.maxWritePacketCount), + callbackState: &interfaceCallbackState{iface: result.iface}, } ReleaseOnCleanup(iface) return iface, nil } func (i *Interface) releaseOnCleanup() { - runtime.AddCleanup(i, func(p unsafe.Pointer) { - C.vmnetRelease(p) - }, objc.Ptr(i)) + runtime.AddCleanup(i, func(state *interfaceCallbackState) { + state.cleanup() + C.vmnetRelease(state.iface) + }, i.callbackState) } // Stop stops I/O on the Interface and releases its associated Network. func (i *Interface) Stop() error { - result := Return(C.VmnetStopInterface(objc.Ptr(i))) - runtime.KeepAlive(i) - if result != ErrSuccess { - return fmt.Errorf("stop vmnet interface: %w", result) + if i == nil || i.callbackState == nil { + return fmt.Errorf("interface is nil") } - return nil + result := i.callbackState.stop() + runtime.KeepAlive(i) + return result } // ReadPackets reads up to packetCount packets into manager. diff --git a/vmnet/vmnet_darwin.h b/vmnet/vmnet_darwin.h index 8a84278..fec1ef7 100644 --- a/vmnet/vmnet_darwin.h +++ b/vmnet/vmnet_darwin.h @@ -40,6 +40,8 @@ void VmnetNetwork_getIPv6Prefix(void *network, struct in6_addr *prefix, uint8_t // MARK: - interface_ref (macOS 26+) +void *VmnetSetPacketsAvailableEventCallback(void *interface, uintptr_t callback, uint32_t *status); +uint32_t VmnetClearPacketsAvailableEventCallback(void *interface, void *queue, uintptr_t callback); uint32_t VmnetStopInterface(void *interface); uint32_t VmnetRead(void *interface, struct vmpktdesc *packets, int *pktcnt); uint32_t VmnetWrite(void *interface, struct vmpktdesc *packets, int *pktcnt); diff --git a/vmnet/vmnet_darwin.m b/vmnet/vmnet_darwin.m index 010b7c4..25b9ee7 100644 --- a/vmnet/vmnet_darwin.m +++ b/vmnet/vmnet_darwin.m @@ -1,4 +1,5 @@ #import "vmnet_darwin.h" +#import // MARK: - CFRelease Wrapper @@ -224,6 +225,48 @@ void VmnetNetwork_getIPv6Prefix(void *network, struct in6_addr *prefix, uint8_t // MARK: - interface_ref (macOS 26+) +extern void callPacketsAvailableEventCallback(uintptr_t handle, int estimatedCount); +extern void releasePacketsAvailableEventCallback(uintptr_t handle); + +void *VmnetSetPacketsAvailableEventCallback(void *iface, uintptr_t callback, uint32_t *status) +{ +#ifdef INCLUDE_TARGET_OSX_26 + if (@available(macOS 26, *)) { + dispatch_queue_t queue = dispatch_queue_create("vmnet.interface.packets", DISPATCH_QUEUE_SERIAL); + *status = vmnet_interface_set_event_callback((interface_ref)iface, VMNET_INTERFACE_PACKETS_AVAILABLE, queue, ^(interface_event_t eventMask, xpc_object_t event) { + if ((eventMask & VMNET_INTERFACE_PACKETS_AVAILABLE) != 0) { + uint64_t estimated = xpc_dictionary_get_uint64(event, vmnet_estimated_packets_available_key); + callPacketsAvailableEventCallback(callback, estimated > INT_MAX ? INT_MAX : (int)estimated); + } + }); + if (*status != VMNET_SUCCESS) { + dispatch_release(queue); + return NULL; + } + return queue; + } +#endif + RAISE_UNSUPPORTED_MACOS_EXCEPTION(); +} + +uint32_t VmnetClearPacketsAvailableEventCallback(void *iface, void *queuePointer, uintptr_t callback) +{ +#ifdef INCLUDE_TARGET_OSX_26 + if (@available(macOS 26, *)) { + vmnet_return_t status = vmnet_interface_set_event_callback((interface_ref)iface, VMNET_INTERFACE_PACKETS_AVAILABLE, NULL, NULL); + if (status == VMNET_SUCCESS) { + dispatch_queue_t queue = (dispatch_queue_t)queuePointer; + dispatch_async(queue, ^{ + releasePacketsAvailableEventCallback(callback); + dispatch_release(queue); + }); + } + return status; + } +#endif + RAISE_UNSUPPORTED_MACOS_EXCEPTION(); +} + uint32_t VmnetRead(void *interface, struct vmpktdesc *packets, int *pktcnt) { #ifdef INCLUDE_TARGET_OSX_26