diff --git a/pkg/aqio/aqio.go b/pkg/aqio/aqio.go index 93256e75..3c3824a1 100644 --- a/pkg/aqio/aqio.go +++ b/pkg/aqio/aqio.go @@ -2,11 +2,16 @@ package aqio import ( "errors" + "fmt" "io" "github.com/johncgriffin/overflow" ) +func NewReadWriteSeeker(buf []byte) *ReadWriteSeeker { + return &ReadWriteSeeker{buf: buf, pos: 0} +} + // ReadWriteSeeker is an in-memory io.ReadWriteSeeker implementation type ReadWriteSeeker struct { buf []byte @@ -15,9 +20,17 @@ type ReadWriteSeeker struct { // Write implements the io.Writer interface func (rws *ReadWriteSeeker) Write(p []byte) (n int, err error) { + fmt.Printf("Write: pos=%d len(p)=%d\n", rws.pos, len(p)) minCap := overflow.Addp(rws.pos, len(p)) if minCap > cap(rws.buf) { // Make sure buf has enough capacity: - buf2 := make([]byte, len(rws.buf), overflow.Addp(minCap, len(p))) // add some extra + fmt.Printf("Write: pos=%d len(p)=%d minCap=%d\n", rws.pos, len(p), minCap) + newCap := cap(rws.buf) * 2 + if newCap == 0 { + newCap = 128 + } else if newCap < minCap { + newCap = minCap * 2 + } + buf2 := make([]byte, len(rws.buf), newCap) // double copy(buf2, rws.buf) rws.buf = buf2 } diff --git a/pkg/c2patypes/stream_adapter.go b/pkg/c2patypes/stream_adapter.go new file mode 100644 index 00000000..c1d8edb0 --- /dev/null +++ b/pkg/c2patypes/stream_adapter.go @@ -0,0 +1,123 @@ +package c2patypes + +import ( + "errors" + "fmt" + "io" +) + +// type Stream interface { +// // Read a stream of bytes from the stream +// ReadStream(length uint64) ([]byte, error) +// // Seek to a position in the stream +// SeekStream(pos int64, mode uint64) (uint64, error) +// // Write a stream of bytes to the stream +// WriteStream(data []byte) (uint64, error) +// } + +// pub enum SeekMode { +// Start = 0, +// End = 1, +// Current = 2, +// } + +const ( + SeekModeStart uint64 = 0 + SeekModeEnd uint64 = 1 + SeekModeCurrent uint64 = 2 +) + +func NewReader(rs io.ReadSeeker) *C2PAStreamReader { + return &C2PAStreamReader{ReadSeeker: rs} +} + +func NewWriter(rws io.ReadWriteSeeker) *C2PAStreamWriter { + return &C2PAStreamWriter{ReadWriteSeeker: rws} +} + +// Wrapped io.ReadSeeker for passing to Rust. Doesn't write. +type C2PAStreamReader struct { + io.ReadSeeker +} + +func (s *C2PAStreamReader) ReadStream(length uint64) ([]byte, error) { + return readStream(s.ReadSeeker, length) +} + +func (s *C2PAStreamReader) SeekStream(pos int64, mode uint64) (uint64, error) { + return seekStream(s.ReadSeeker, pos, mode) +} + +func (s *C2PAStreamReader) WriteStream(data []byte) (uint64, error) { + return 0, fmt.Errorf("Writing is not implemented for C2PAStreamReader") +} + +// Wrapped io.Writer for passing to Rust. +type C2PAStreamWriter struct { + io.ReadWriteSeeker +} + +func (s *C2PAStreamWriter) ReadStream(length uint64) ([]byte, error) { + return readStream(s.ReadWriteSeeker, length) +} + +func (s *C2PAStreamWriter) SeekStream(pos int64, mode uint64) (uint64, error) { + return seekStream(s.ReadWriteSeeker, pos, mode) +} + +func (s *C2PAStreamWriter) WriteStream(data []byte) (uint64, error) { + return writeStream(s.ReadWriteSeeker, data) +} + +func readStream(r io.ReadSeeker, length uint64) ([]byte, error) { + // fmt.Printf("read length=%d\n", length) + bs := make([]byte, length) + read, err := r.Read(bs) + if err != nil { + if errors.Is(err, io.EOF) { + if read == 0 { + // fmt.Printf("read EOF read=%d returning empty?", read) + return []byte{}, nil + } + // partial := bs[read:] + // return partial, nil + } + // fmt.Printf("io error=%s\n", err) + return []byte{}, err + } + if uint64(read) < length { + partial := bs[:read] + // fmt.Printf("read returning partial read=%d len=%d\n", read, len(partial)) + return partial, nil + } + // fmt.Printf("read returning full read=%d len=%d\n", read, len(bs)) + return bs, nil +} + +func seekStream(r io.ReadSeeker, pos int64, mode uint64) (uint64, error) { + // fmt.Printf("seek pos=%d\n", pos) + var seekMode int + if mode == SeekModeCurrent { + seekMode = io.SeekCurrent + } else if mode == SeekModeStart { + seekMode = io.SeekStart + } else if mode == SeekModeEnd { + seekMode = io.SeekEnd + } else { + // fmt.Printf("seek mode unsupported mode=%d\n", mode) + return 0, fmt.Errorf("unknown seek mode: %d", mode) + } + newPos, err := r.Seek(pos, seekMode) + if err != nil { + return 0, err + } + return uint64(newPos), nil +} + +func writeStream(w io.ReadWriteSeeker, data []byte) (uint64, error) { + wrote, err := w.Write(data) + if err != nil { + return uint64(wrote), err + } + return uint64(wrote), nil +} diff --git a/pkg/iroh/generated/iroh_streamplace/iroh_streamplace.go b/pkg/iroh/generated/iroh_streamplace/iroh_streamplace.go index 5d85520a..8312662f 100644 --- a/pkg/iroh/generated/iroh_streamplace/iroh_streamplace.go +++ b/pkg/iroh/generated/iroh_streamplace/iroh_streamplace.go @@ -338,6 +338,7 @@ func init() { FfiConverterDataHandlerINSTANCE.register() FfiConverterGoSignerINSTANCE.register() + FfiConverterStreamINSTANCE.register() uniffiCheckChecksums() } @@ -356,7 +357,7 @@ func uniffiCheckChecksums() { checksum := rustCall(func(_uniffiStatus *C.RustCallStatus) C.uint16_t { return C.uniffi_iroh_streamplace_checksum_func_get_manifest_and_cert() }) - if checksum != 17550 { + if checksum != 36028 { // If this happens try cleaning and rebuilding your project panic("iroh_streamplace: uniffi_iroh_streamplace_checksum_func_get_manifest_and_cert: UniFFI API checksum mismatch") } @@ -365,7 +366,7 @@ func uniffiCheckChecksums() { checksum := rustCall(func(_uniffiStatus *C.RustCallStatus) C.uint16_t { return C.uniffi_iroh_streamplace_checksum_func_get_manifests() }) - if checksum != 17 { + if checksum != 2548 { // If this happens try cleaning and rebuilding your project panic("iroh_streamplace: uniffi_iroh_streamplace_checksum_func_get_manifests: UniFFI API checksum mismatch") } @@ -401,7 +402,7 @@ func uniffiCheckChecksums() { checksum := rustCall(func(_uniffiStatus *C.RustCallStatus) C.uint16_t { return C.uniffi_iroh_streamplace_checksum_func_resign() }) - if checksum != 32728 { + if checksum != 7588 { // If this happens try cleaning and rebuilding your project panic("iroh_streamplace: uniffi_iroh_streamplace_checksum_func_resign: UniFFI API checksum mismatch") } @@ -410,7 +411,7 @@ func uniffiCheckChecksums() { checksum := rustCall(func(_uniffiStatus *C.RustCallStatus) C.uint16_t { return C.uniffi_iroh_streamplace_checksum_func_sign() }) - if checksum != 23786 { + if checksum != 50601 { // If this happens try cleaning and rebuilding your project panic("iroh_streamplace: uniffi_iroh_streamplace_checksum_func_sign: UniFFI API checksum mismatch") } @@ -419,7 +420,7 @@ func uniffiCheckChecksums() { checksum := rustCall(func(_uniffiStatus *C.RustCallStatus) C.uint16_t { return C.uniffi_iroh_streamplace_checksum_func_sign_with_ingredients() }) - if checksum != 63680 { + if checksum != 34840 { // If this happens try cleaning and rebuilding your project panic("iroh_streamplace: uniffi_iroh_streamplace_checksum_func_sign_with_ingredients: UniFFI API checksum mismatch") } @@ -712,6 +713,33 @@ func uniffiCheckChecksums() { panic("iroh_streamplace: uniffi_iroh_streamplace_checksum_method_publickey_fmt_short: UniFFI API checksum mismatch") } } + { + checksum := rustCall(func(_uniffiStatus *C.RustCallStatus) C.uint16_t { + return C.uniffi_iroh_streamplace_checksum_method_stream_read_stream() + }) + if checksum != 62815 { + // If this happens try cleaning and rebuilding your project + panic("iroh_streamplace: uniffi_iroh_streamplace_checksum_method_stream_read_stream: UniFFI API checksum mismatch") + } + } + { + checksum := rustCall(func(_uniffiStatus *C.RustCallStatus) C.uint16_t { + return C.uniffi_iroh_streamplace_checksum_method_stream_seek_stream() + }) + if checksum != 56397 { + // If this happens try cleaning and rebuilding your project + panic("iroh_streamplace: uniffi_iroh_streamplace_checksum_method_stream_seek_stream: UniFFI API checksum mismatch") + } + } + { + checksum := rustCall(func(_uniffiStatus *C.RustCallStatus) C.uint16_t { + return C.uniffi_iroh_streamplace_checksum_method_stream_write_stream() + }) + if checksum != 59847 { + // If this happens try cleaning and rebuilding your project + panic("iroh_streamplace: uniffi_iroh_streamplace_checksum_method_stream_write_stream: UniFFI API checksum mismatch") + } + } { checksum := rustCall(func(_uniffiStatus *C.RustCallStatus) C.uint16_t { return C.uniffi_iroh_streamplace_checksum_method_subscriberesponse_next_raw() @@ -819,6 +847,30 @@ type FfiDestroyerUint64 struct{} func (FfiDestroyerUint64) Destroy(_ uint64) {} +type FfiConverterInt64 struct{} + +var FfiConverterInt64INSTANCE = FfiConverterInt64{} + +func (FfiConverterInt64) Lower(value int64) C.int64_t { + return C.int64_t(value) +} + +func (FfiConverterInt64) Write(writer io.Writer, value int64) { + writeInt64(writer, value) +} + +func (FfiConverterInt64) Lift(value C.int64_t) int64 { + return int64(value) +} + +func (FfiConverterInt64) Read(reader io.Reader) int64 { + return readInt64(reader) +} + +type FfiDestroyerInt64 struct{} + +func (FfiDestroyerInt64) Destroy(_ int64) {} + type FfiConverterBool struct{} var FfiConverterBoolINSTANCE = FfiConverterBool{} @@ -2393,6 +2445,247 @@ func (_ FfiDestroyerPublicKey) Destroy(value *PublicKey) { value.Destroy() } +// This allows for a callback stream over the Uniffi interface. +// Implement these stream functions in the foreign language +// and this will provide Rust Stream trait implementations +// This is necessary since the Rust traits cannot be implemented directly +// as uniffi callbacks +type Stream interface { + // Read a stream of bytes from the stream + ReadStream(length uint64) ([]byte, error) + // Seek to a position in the stream + SeekStream(pos int64, mode uint64) (uint64, error) + // Write a stream of bytes to the stream + WriteStream(data []byte) (uint64, error) +} + +// This allows for a callback stream over the Uniffi interface. +// Implement these stream functions in the foreign language +// and this will provide Rust Stream trait implementations +// This is necessary since the Rust traits cannot be implemented directly +// as uniffi callbacks +type StreamImpl struct { + ffiObject FfiObject +} + +// Read a stream of bytes from the stream +func (_self *StreamImpl) ReadStream(length uint64) ([]byte, error) { + _pointer := _self.ffiObject.incrementPointer("Stream") + defer _self.ffiObject.decrementPointer() + _uniffiRV, _uniffiErr := rustCallWithError[SpError](FfiConverterSpError{}, func(_uniffiStatus *C.RustCallStatus) RustBufferI { + return GoRustBuffer{ + inner: C.uniffi_iroh_streamplace_fn_method_stream_read_stream( + _pointer, FfiConverterUint64INSTANCE.Lower(length), _uniffiStatus), + } + }) + if _uniffiErr != nil { + var _uniffiDefaultValue []byte + return _uniffiDefaultValue, _uniffiErr + } else { + return FfiConverterBytesINSTANCE.Lift(_uniffiRV), nil + } +} + +// Seek to a position in the stream +func (_self *StreamImpl) SeekStream(pos int64, mode uint64) (uint64, error) { + _pointer := _self.ffiObject.incrementPointer("Stream") + defer _self.ffiObject.decrementPointer() + _uniffiRV, _uniffiErr := rustCallWithError[SpError](FfiConverterSpError{}, func(_uniffiStatus *C.RustCallStatus) C.uint64_t { + return C.uniffi_iroh_streamplace_fn_method_stream_seek_stream( + _pointer, FfiConverterInt64INSTANCE.Lower(pos), FfiConverterUint64INSTANCE.Lower(mode), _uniffiStatus) + }) + if _uniffiErr != nil { + var _uniffiDefaultValue uint64 + return _uniffiDefaultValue, _uniffiErr + } else { + return FfiConverterUint64INSTANCE.Lift(_uniffiRV), nil + } +} + +// Write a stream of bytes to the stream +func (_self *StreamImpl) WriteStream(data []byte) (uint64, error) { + _pointer := _self.ffiObject.incrementPointer("Stream") + defer _self.ffiObject.decrementPointer() + _uniffiRV, _uniffiErr := rustCallWithError[SpError](FfiConverterSpError{}, func(_uniffiStatus *C.RustCallStatus) C.uint64_t { + return C.uniffi_iroh_streamplace_fn_method_stream_write_stream( + _pointer, FfiConverterBytesINSTANCE.Lower(data), _uniffiStatus) + }) + if _uniffiErr != nil { + var _uniffiDefaultValue uint64 + return _uniffiDefaultValue, _uniffiErr + } else { + return FfiConverterUint64INSTANCE.Lift(_uniffiRV), nil + } +} +func (object *StreamImpl) Destroy() { + runtime.SetFinalizer(object, nil) + object.ffiObject.destroy() +} + +type FfiConverterStream struct { + handleMap *concurrentHandleMap[Stream] +} + +var FfiConverterStreamINSTANCE = FfiConverterStream{ + handleMap: newConcurrentHandleMap[Stream](), +} + +func (c FfiConverterStream) Lift(pointer unsafe.Pointer) Stream { + result := &StreamImpl{ + newFfiObject( + pointer, + func(pointer unsafe.Pointer, status *C.RustCallStatus) unsafe.Pointer { + return C.uniffi_iroh_streamplace_fn_clone_stream(pointer, status) + }, + func(pointer unsafe.Pointer, status *C.RustCallStatus) { + C.uniffi_iroh_streamplace_fn_free_stream(pointer, status) + }, + ), + } + runtime.SetFinalizer(result, (*StreamImpl).Destroy) + return result +} + +func (c FfiConverterStream) Read(reader io.Reader) Stream { + return c.Lift(unsafe.Pointer(uintptr(readUint64(reader)))) +} + +func (c FfiConverterStream) Lower(value Stream) unsafe.Pointer { + // TODO: this is bad - all synchronization from ObjectRuntime.go is discarded here, + // because the pointer will be decremented immediately after this function returns, + // and someone will be left holding onto a non-locked pointer. + pointer := unsafe.Pointer(uintptr(c.handleMap.insert(value))) + return pointer + +} + +func (c FfiConverterStream) Write(writer io.Writer, value Stream) { + writeUint64(writer, uint64(uintptr(c.Lower(value)))) +} + +type FfiDestroyerStream struct{} + +func (_ FfiDestroyerStream) Destroy(value Stream) { + if val, ok := value.(*StreamImpl); ok { + val.Destroy() + } else { + panic("Expected *StreamImpl") + } +} + +//export iroh_streamplace_cgo_dispatchCallbackInterfaceStreamMethod0 +func iroh_streamplace_cgo_dispatchCallbackInterfaceStreamMethod0(uniffiHandle C.uint64_t, length C.uint64_t, uniffiOutReturn *C.RustBuffer, callStatus *C.RustCallStatus) { + handle := uint64(uniffiHandle) + uniffiObj, ok := FfiConverterStreamINSTANCE.handleMap.tryGet(handle) + if !ok { + panic(fmt.Errorf("no callback in handle map: %d", handle)) + } + + res, err := + uniffiObj.ReadStream( + FfiConverterUint64INSTANCE.Lift(length), + ) + + if err != nil { + var actualError *SpError + if errors.As(err, &actualError) { + *callStatus = C.RustCallStatus{ + code: C.int8_t(uniffiCallbackResultError), + errorBuf: FfiConverterSpErrorINSTANCE.Lower(actualError), + } + } else { + *callStatus = C.RustCallStatus{ + code: C.int8_t(uniffiCallbackUnexpectedResultError), + } + } + return + } + + *uniffiOutReturn = FfiConverterBytesINSTANCE.Lower(res) +} + +//export iroh_streamplace_cgo_dispatchCallbackInterfaceStreamMethod1 +func iroh_streamplace_cgo_dispatchCallbackInterfaceStreamMethod1(uniffiHandle C.uint64_t, pos C.int64_t, mode C.uint64_t, uniffiOutReturn *C.uint64_t, callStatus *C.RustCallStatus) { + handle := uint64(uniffiHandle) + uniffiObj, ok := FfiConverterStreamINSTANCE.handleMap.tryGet(handle) + if !ok { + panic(fmt.Errorf("no callback in handle map: %d", handle)) + } + + res, err := + uniffiObj.SeekStream( + FfiConverterInt64INSTANCE.Lift(pos), + FfiConverterUint64INSTANCE.Lift(mode), + ) + + if err != nil { + var actualError *SpError + if errors.As(err, &actualError) { + *callStatus = C.RustCallStatus{ + code: C.int8_t(uniffiCallbackResultError), + errorBuf: FfiConverterSpErrorINSTANCE.Lower(actualError), + } + } else { + *callStatus = C.RustCallStatus{ + code: C.int8_t(uniffiCallbackUnexpectedResultError), + } + } + return + } + + *uniffiOutReturn = FfiConverterUint64INSTANCE.Lower(res) +} + +//export iroh_streamplace_cgo_dispatchCallbackInterfaceStreamMethod2 +func iroh_streamplace_cgo_dispatchCallbackInterfaceStreamMethod2(uniffiHandle C.uint64_t, data C.RustBuffer, uniffiOutReturn *C.uint64_t, callStatus *C.RustCallStatus) { + handle := uint64(uniffiHandle) + uniffiObj, ok := FfiConverterStreamINSTANCE.handleMap.tryGet(handle) + if !ok { + panic(fmt.Errorf("no callback in handle map: %d", handle)) + } + + res, err := + uniffiObj.WriteStream( + FfiConverterBytesINSTANCE.Lift(GoRustBuffer{ + inner: data, + }), + ) + + if err != nil { + var actualError *SpError + if errors.As(err, &actualError) { + *callStatus = C.RustCallStatus{ + code: C.int8_t(uniffiCallbackResultError), + errorBuf: FfiConverterSpErrorINSTANCE.Lower(actualError), + } + } else { + *callStatus = C.RustCallStatus{ + code: C.int8_t(uniffiCallbackUnexpectedResultError), + } + } + return + } + + *uniffiOutReturn = FfiConverterUint64INSTANCE.Lower(res) +} + +var UniffiVTableCallbackInterfaceStreamINSTANCE = C.UniffiVTableCallbackInterfaceStream{ + readStream: (C.UniffiCallbackInterfaceStreamMethod0)(C.iroh_streamplace_cgo_dispatchCallbackInterfaceStreamMethod0), + seekStream: (C.UniffiCallbackInterfaceStreamMethod1)(C.iroh_streamplace_cgo_dispatchCallbackInterfaceStreamMethod1), + writeStream: (C.UniffiCallbackInterfaceStreamMethod2)(C.iroh_streamplace_cgo_dispatchCallbackInterfaceStreamMethod2), + + uniffiFree: (C.UniffiCallbackInterfaceFree)(C.iroh_streamplace_cgo_dispatchCallbackInterfaceStreamFree), +} + +//export iroh_streamplace_cgo_dispatchCallbackInterfaceStreamFree +func iroh_streamplace_cgo_dispatchCallbackInterfaceStreamFree(handle C.uint64_t) { + FfiConverterStreamINSTANCE.handleMap.remove(uint64(handle)) +} + +func (c FfiConverterStream) register() { + C.uniffi_iroh_streamplace_fn_init_callback_vtable_stream(&UniffiVTableCallbackInterfaceStreamINSTANCE) +} + // A response to a subscribe request. // // This can be used as a stream of [`SubscribeItem`]s. @@ -3613,6 +3906,7 @@ func (err SpError) Unwrap() error { // Err* are used for checking error type with `errors.Is` var ErrSpErrorNoCertificateChainFound = fmt.Errorf("SpErrorNoCertificateChainFound") var ErrSpErrorC2paError = fmt.Errorf("SpErrorC2paError") +var ErrSpErrorIoError = fmt.Errorf("SpErrorIoError") // Variant structs type SpErrorNoCertificateChainFound struct { @@ -3653,6 +3947,25 @@ func (self SpErrorC2paError) Is(target error) bool { return target == ErrSpErrorC2paError } +type SpErrorIoError struct { + message string +} + +func NewSpErrorIoError() *SpError { + return &SpError{err: &SpErrorIoError{}} +} + +func (e SpErrorIoError) destroy() { +} + +func (err SpErrorIoError) Error() string { + return fmt.Sprintf("IoError: %s", err.message) +} + +func (self SpErrorIoError) Is(target error) bool { + return target == ErrSpErrorIoError +} + type FfiConverterSpError struct{} var FfiConverterSpErrorINSTANCE = FfiConverterSpError{} @@ -3674,6 +3987,8 @@ func (c FfiConverterSpError) Read(reader io.Reader) *SpError { return &SpError{&SpErrorNoCertificateChainFound{message}} case 2: return &SpError{&SpErrorC2paError{message}} + case 3: + return &SpError{&SpErrorIoError{message}} default: panic(fmt.Sprintf("Unknown error code %d in FfiConverterSpError.Read()", errorID)) } @@ -3686,6 +4001,8 @@ func (c FfiConverterSpError) Write(writer io.Writer, value *SpError) { writeInt32(writer, 1) case *SpErrorC2paError: writeInt32(writer, 2) + case *SpErrorIoError: + writeInt32(writer, 3) default: _ = variantValue panic(fmt.Sprintf("invalid error value `%v` in FfiConverterSpError.Write", value)) @@ -3700,6 +4017,8 @@ func (_ FfiDestroyerSpError) Destroy(value *SpError) { variantValue.destroy() case SpErrorC2paError: variantValue.destroy() + case SpErrorIoError: + variantValue.destroy() default: _ = variantValue panic(fmt.Sprintf("invalid error value `%v` in FfiDestroyerSpError.Destroy", value)) @@ -4759,10 +5078,10 @@ func iroh_streamplace_uniffiFreeGorutine(data C.uint64_t) { guard <- struct{}{} } -func GetManifestAndCert(data []byte) (string, error) { +func GetManifestAndCert(data Stream) (string, error) { _uniffiRV, _uniffiErr := rustCallWithError[SpError](FfiConverterSpError{}, func(_uniffiStatus *C.RustCallStatus) RustBufferI { return GoRustBuffer{ - inner: C.uniffi_iroh_streamplace_fn_func_get_manifest_and_cert(FfiConverterBytesINSTANCE.Lower(data), _uniffiStatus), + inner: C.uniffi_iroh_streamplace_fn_func_get_manifest_and_cert(FfiConverterStreamINSTANCE.Lower(data), _uniffiStatus), } }) if _uniffiErr != nil { @@ -4773,10 +5092,10 @@ func GetManifestAndCert(data []byte) (string, error) { } } -func GetManifests(data []byte) (string, error) { +func GetManifests(data Stream) (string, error) { _uniffiRV, _uniffiErr := rustCallWithError[SpError](FfiConverterSpError{}, func(_uniffiStatus *C.RustCallStatus) RustBufferI { return GoRustBuffer{ - inner: C.uniffi_iroh_streamplace_fn_func_get_manifests(FfiConverterBytesINSTANCE.Lower(data), _uniffiStatus), + inner: C.uniffi_iroh_streamplace_fn_func_get_manifests(FfiConverterStreamINSTANCE.Lower(data), _uniffiStatus), } }) if _uniffiErr != nil { @@ -4821,10 +5140,10 @@ func NodeIdFromTicket(ticketStr string) (*PublicKey, error) { } } -func Resign(unsignedSegData [][]byte, signedConcatData []byte, manifestList []string, certs [][]byte) ([][]byte, error) { +func Resign(unsignedSegData [][]byte, signedConcatData Stream, manifestList []string, certs [][]byte) ([][]byte, error) { _uniffiRV, _uniffiErr := rustCallWithError[SpError](FfiConverterSpError{}, func(_uniffiStatus *C.RustCallStatus) RustBufferI { return GoRustBuffer{ - inner: C.uniffi_iroh_streamplace_fn_func_resign(FfiConverterSequenceBytesINSTANCE.Lower(unsignedSegData), FfiConverterBytesINSTANCE.Lower(signedConcatData), FfiConverterSequenceStringINSTANCE.Lower(manifestList), FfiConverterSequenceBytesINSTANCE.Lower(certs), _uniffiStatus), + inner: C.uniffi_iroh_streamplace_fn_func_resign(FfiConverterSequenceBytesINSTANCE.Lower(unsignedSegData), FfiConverterStreamINSTANCE.Lower(signedConcatData), FfiConverterSequenceStringINSTANCE.Lower(manifestList), FfiConverterSequenceBytesINSTANCE.Lower(certs), _uniffiStatus), } }) if _uniffiErr != nil { @@ -4835,10 +5154,10 @@ func Resign(unsignedSegData [][]byte, signedConcatData []byte, manifestList []st } } -func Sign(manifest string, data []byte, certs []byte, gosigner GoSigner) ([]byte, error) { +func Sign(manifest string, data Stream, certs []byte, gosigner GoSigner) ([]byte, error) { _uniffiRV, _uniffiErr := rustCallWithError[SpError](FfiConverterSpError{}, func(_uniffiStatus *C.RustCallStatus) RustBufferI { return GoRustBuffer{ - inner: C.uniffi_iroh_streamplace_fn_func_sign(FfiConverterStringINSTANCE.Lower(manifest), FfiConverterBytesINSTANCE.Lower(data), FfiConverterBytesINSTANCE.Lower(certs), FfiConverterGoSignerINSTANCE.Lower(gosigner), _uniffiStatus), + inner: C.uniffi_iroh_streamplace_fn_func_sign(FfiConverterStringINSTANCE.Lower(manifest), FfiConverterStreamINSTANCE.Lower(data), FfiConverterBytesINSTANCE.Lower(certs), FfiConverterGoSignerINSTANCE.Lower(gosigner), _uniffiStatus), } }) if _uniffiErr != nil { @@ -4849,18 +5168,12 @@ func Sign(manifest string, data []byte, certs []byte, gosigner GoSigner) ([]byte } } -func SignWithIngredients(manifest string, data []byte, certs []byte, ingredients [][]byte, gosigner GoSigner) ([]byte, error) { - _uniffiRV, _uniffiErr := rustCallWithError[SpError](FfiConverterSpError{}, func(_uniffiStatus *C.RustCallStatus) RustBufferI { - return GoRustBuffer{ - inner: C.uniffi_iroh_streamplace_fn_func_sign_with_ingredients(FfiConverterStringINSTANCE.Lower(manifest), FfiConverterBytesINSTANCE.Lower(data), FfiConverterBytesINSTANCE.Lower(certs), FfiConverterSequenceBytesINSTANCE.Lower(ingredients), FfiConverterGoSignerINSTANCE.Lower(gosigner), _uniffiStatus), - } +func SignWithIngredients(manifest string, data Stream, certs []byte, ingredients [][]byte, gosigner GoSigner, output Stream) error { + _, _uniffiErr := rustCallWithError[SpError](FfiConverterSpError{}, func(_uniffiStatus *C.RustCallStatus) bool { + C.uniffi_iroh_streamplace_fn_func_sign_with_ingredients(FfiConverterStringINSTANCE.Lower(manifest), FfiConverterStreamINSTANCE.Lower(data), FfiConverterBytesINSTANCE.Lower(certs), FfiConverterSequenceBytesINSTANCE.Lower(ingredients), FfiConverterGoSignerINSTANCE.Lower(gosigner), FfiConverterStreamINSTANCE.Lower(output), _uniffiStatus) + return false }) - if _uniffiErr != nil { - var _uniffiDefaultValue []byte - return _uniffiDefaultValue, _uniffiErr - } else { - return FfiConverterBytesINSTANCE.Lift(_uniffiRV), nil - } + return _uniffiErr.AsError() } func SubscribeItemDebug(item SubscribeItem) string { diff --git a/pkg/iroh/generated/iroh_streamplace/iroh_streamplace.h b/pkg/iroh/generated/iroh_streamplace/iroh_streamplace.h index 0012483d..e41afb4d 100644 --- a/pkg/iroh/generated/iroh_streamplace/iroh_streamplace.h +++ b/pkg/iroh/generated/iroh_streamplace/iroh_streamplace.h @@ -405,6 +405,48 @@ static void call_UniffiCallbackInterfaceGoSignerMethod0( } +#endif +#ifndef UNIFFI_FFIDEF_CALLBACK_INTERFACE_STREAM_METHOD0 +#define UNIFFI_FFIDEF_CALLBACK_INTERFACE_STREAM_METHOD0 +typedef void (*UniffiCallbackInterfaceStreamMethod0)(uint64_t uniffi_handle, uint64_t length, RustBuffer* uniffi_out_return, RustCallStatus* callStatus ); + +// Making function static works arround: +// https://github.com/golang/go/issues/11263 +static void call_UniffiCallbackInterfaceStreamMethod0( + UniffiCallbackInterfaceStreamMethod0 cb, uint64_t uniffi_handle, uint64_t length, RustBuffer* uniffi_out_return, RustCallStatus* callStatus ) +{ + return cb(uniffi_handle, length, uniffi_out_return, callStatus ); +} + + +#endif +#ifndef UNIFFI_FFIDEF_CALLBACK_INTERFACE_STREAM_METHOD1 +#define UNIFFI_FFIDEF_CALLBACK_INTERFACE_STREAM_METHOD1 +typedef void (*UniffiCallbackInterfaceStreamMethod1)(uint64_t uniffi_handle, int64_t pos, uint64_t mode, uint64_t* uniffi_out_return, RustCallStatus* callStatus ); + +// Making function static works arround: +// https://github.com/golang/go/issues/11263 +static void call_UniffiCallbackInterfaceStreamMethod1( + UniffiCallbackInterfaceStreamMethod1 cb, uint64_t uniffi_handle, int64_t pos, uint64_t mode, uint64_t* uniffi_out_return, RustCallStatus* callStatus ) +{ + return cb(uniffi_handle, pos, mode, uniffi_out_return, callStatus ); +} + + +#endif +#ifndef UNIFFI_FFIDEF_CALLBACK_INTERFACE_STREAM_METHOD2 +#define UNIFFI_FFIDEF_CALLBACK_INTERFACE_STREAM_METHOD2 +typedef void (*UniffiCallbackInterfaceStreamMethod2)(uint64_t uniffi_handle, RustBuffer data, uint64_t* uniffi_out_return, RustCallStatus* callStatus ); + +// Making function static works arround: +// https://github.com/golang/go/issues/11263 +static void call_UniffiCallbackInterfaceStreamMethod2( + UniffiCallbackInterfaceStreamMethod2 cb, uint64_t uniffi_handle, RustBuffer data, uint64_t* uniffi_out_return, RustCallStatus* callStatus ) +{ + return cb(uniffi_handle, data, uniffi_out_return, callStatus ); +} + + #endif #ifndef UNIFFI_FFIDEF_V_TABLE_CALLBACK_INTERFACE_DATA_HANDLER #define UNIFFI_FFIDEF_V_TABLE_CALLBACK_INTERFACE_DATA_HANDLER @@ -421,6 +463,16 @@ typedef struct UniffiVTableCallbackInterfaceGoSigner { UniffiCallbackInterfaceFree uniffiFree; } UniffiVTableCallbackInterfaceGoSigner; +#endif +#ifndef UNIFFI_FFIDEF_V_TABLE_CALLBACK_INTERFACE_STREAM +#define UNIFFI_FFIDEF_V_TABLE_CALLBACK_INTERFACE_STREAM +typedef struct UniffiVTableCallbackInterfaceStream { + UniffiCallbackInterfaceStreamMethod0 readStream; + UniffiCallbackInterfaceStreamMethod1 seekStream; + UniffiCallbackInterfaceStreamMethod2 writeStream; + UniffiCallbackInterfaceFree uniffiFree; +} UniffiVTableCallbackInterfaceStream; + #endif #ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_CLONE_DATAHANDLER #define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_CLONE_DATAHANDLER @@ -698,6 +750,36 @@ RustBuffer uniffi_iroh_streamplace_fn_method_publickey_fmt_short(void* ptr, Rust RustBuffer uniffi_iroh_streamplace_fn_method_publickey_uniffi_trait_display(void* ptr, RustCallStatus *out_status ); #endif +#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_CLONE_STREAM +#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_CLONE_STREAM +void* uniffi_iroh_streamplace_fn_clone_stream(void* ptr, RustCallStatus *out_status +); +#endif +#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_FREE_STREAM +#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_FREE_STREAM +void uniffi_iroh_streamplace_fn_free_stream(void* ptr, RustCallStatus *out_status +); +#endif +#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_INIT_CALLBACK_VTABLE_STREAM +#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_INIT_CALLBACK_VTABLE_STREAM +void uniffi_iroh_streamplace_fn_init_callback_vtable_stream(UniffiVTableCallbackInterfaceStream* vtable +); +#endif +#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_METHOD_STREAM_READ_STREAM +#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_METHOD_STREAM_READ_STREAM +RustBuffer uniffi_iroh_streamplace_fn_method_stream_read_stream(void* ptr, uint64_t length, RustCallStatus *out_status +); +#endif +#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_METHOD_STREAM_SEEK_STREAM +#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_METHOD_STREAM_SEEK_STREAM +uint64_t uniffi_iroh_streamplace_fn_method_stream_seek_stream(void* ptr, int64_t pos, uint64_t mode, RustCallStatus *out_status +); +#endif +#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_METHOD_STREAM_WRITE_STREAM +#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_METHOD_STREAM_WRITE_STREAM +uint64_t uniffi_iroh_streamplace_fn_method_stream_write_stream(void* ptr, RustBuffer data, RustCallStatus *out_status +); +#endif #ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_CLONE_SUBSCRIBERESPONSE #define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_CLONE_SUBSCRIBERESPONSE void* uniffi_iroh_streamplace_fn_clone_subscriberesponse(void* ptr, RustCallStatus *out_status @@ -735,12 +817,12 @@ uint64_t uniffi_iroh_streamplace_fn_method_writescope_put(void* ptr, RustBuffer #endif #ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_FUNC_GET_MANIFEST_AND_CERT #define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_FUNC_GET_MANIFEST_AND_CERT -RustBuffer uniffi_iroh_streamplace_fn_func_get_manifest_and_cert(RustBuffer data, RustCallStatus *out_status +RustBuffer uniffi_iroh_streamplace_fn_func_get_manifest_and_cert(void* data, RustCallStatus *out_status ); #endif #ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_FUNC_GET_MANIFESTS #define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_FUNC_GET_MANIFESTS -RustBuffer uniffi_iroh_streamplace_fn_func_get_manifests(RustBuffer data, RustCallStatus *out_status +RustBuffer uniffi_iroh_streamplace_fn_func_get_manifests(void* data, RustCallStatus *out_status ); #endif #ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_FUNC_INIT_LOGGING @@ -761,17 +843,17 @@ void* uniffi_iroh_streamplace_fn_func_node_id_from_ticket(RustBuffer ticket_str, #endif #ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_FUNC_RESIGN #define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_FUNC_RESIGN -RustBuffer uniffi_iroh_streamplace_fn_func_resign(RustBuffer unsigned_seg_data, RustBuffer signed_concat_data, RustBuffer manifest_list, RustBuffer certs, RustCallStatus *out_status +RustBuffer uniffi_iroh_streamplace_fn_func_resign(RustBuffer unsigned_seg_data, void* signed_concat_data, RustBuffer manifest_list, RustBuffer certs, RustCallStatus *out_status ); #endif #ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_FUNC_SIGN #define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_FUNC_SIGN -RustBuffer uniffi_iroh_streamplace_fn_func_sign(RustBuffer manifest, RustBuffer data, RustBuffer certs, void* gosigner, RustCallStatus *out_status +RustBuffer uniffi_iroh_streamplace_fn_func_sign(RustBuffer manifest, void* data, RustBuffer certs, void* gosigner, RustCallStatus *out_status ); #endif #ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_FUNC_SIGN_WITH_INGREDIENTS #define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_FUNC_SIGN_WITH_INGREDIENTS -RustBuffer uniffi_iroh_streamplace_fn_func_sign_with_ingredients(RustBuffer manifest, RustBuffer data, RustBuffer certs, RustBuffer ingredients, void* gosigner, RustCallStatus *out_status +void uniffi_iroh_streamplace_fn_func_sign_with_ingredients(RustBuffer manifest, void* data, RustBuffer certs, RustBuffer ingredients, void* gosigner, void* output, RustCallStatus *out_status ); #endif #ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_FN_FUNC_SUBSCRIBE_ITEM_DEBUG @@ -1297,6 +1379,24 @@ uint16_t uniffi_iroh_streamplace_checksum_method_publickey_equal(void #define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_METHOD_PUBLICKEY_FMT_SHORT uint16_t uniffi_iroh_streamplace_checksum_method_publickey_fmt_short(void +); +#endif +#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_METHOD_STREAM_READ_STREAM +#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_METHOD_STREAM_READ_STREAM +uint16_t uniffi_iroh_streamplace_checksum_method_stream_read_stream(void + +); +#endif +#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_METHOD_STREAM_SEEK_STREAM +#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_METHOD_STREAM_SEEK_STREAM +uint16_t uniffi_iroh_streamplace_checksum_method_stream_seek_stream(void + +); +#endif +#ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_METHOD_STREAM_WRITE_STREAM +#define UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_METHOD_STREAM_WRITE_STREAM +uint16_t uniffi_iroh_streamplace_checksum_method_stream_write_stream(void + ); #endif #ifndef UNIFFI_FFIDEF_UNIFFI_IROH_STREAMPLACE_CHECKSUM_METHOD_SUBSCRIBERESPONSE_NEXT_RAW @@ -1364,6 +1464,10 @@ uint32_t ffi_iroh_streamplace_uniffi_contract_version(void void iroh_streamplace_cgo_dispatchCallbackInterfaceDataHandlerFree(uint64_t handle); void iroh_streamplace_cgo_dispatchCallbackInterfaceGoSignerMethod0(uint64_t uniffi_handle, RustBuffer data, RustBuffer* uniffi_out_return, RustCallStatus* callStatus ); void iroh_streamplace_cgo_dispatchCallbackInterfaceGoSignerFree(uint64_t handle); + void iroh_streamplace_cgo_dispatchCallbackInterfaceStreamMethod0(uint64_t uniffi_handle, uint64_t length, RustBuffer* uniffi_out_return, RustCallStatus* callStatus ); + void iroh_streamplace_cgo_dispatchCallbackInterfaceStreamMethod1(uint64_t uniffi_handle, int64_t pos, uint64_t mode, uint64_t* uniffi_out_return, RustCallStatus* callStatus ); + void iroh_streamplace_cgo_dispatchCallbackInterfaceStreamMethod2(uint64_t uniffi_handle, RustBuffer data, uint64_t* uniffi_out_return, RustCallStatus* callStatus ); + void iroh_streamplace_cgo_dispatchCallbackInterfaceStreamFree(uint64_t handle); void iroh_streamplace_uniffiFutureContinuationCallback(uint64_t, int8_t); void iroh_streamplace_uniffiFreeGorutine(uint64_t); diff --git a/pkg/media/media_signer.go b/pkg/media/media_signer.go index 24dedd96..3396fdee 100644 --- a/pkg/media/media_signer.go +++ b/pkg/media/media_signer.go @@ -12,6 +12,7 @@ import ( "time" "go.opentelemetry.io/otel" + "stream.place/streamplace/pkg/aqio" "stream.place/streamplace/pkg/aqtime" "stream.place/streamplace/pkg/atproto" c2patypes "stream.place/streamplace/pkg/c2patypes" @@ -159,7 +160,7 @@ func (ms *MediaSignerLocal) SignMP4(ctx context.Context, input io.ReadSeeker, st rustCallbackSigner := &RustCallbackSigner{ Signer: ms.Signer, } - bs, err = iroh_streamplace.Sign(string(manifestBs), bs, ms.Cert, rustCallbackSigner) + bs, err = iroh_streamplace.Sign(string(manifestBs), c2patypes.NewReader(aqio.NewReadWriteSeeker(bs)), ms.Cert, rustCallbackSigner) if err != nil { return nil, err } @@ -179,12 +180,12 @@ func (ms *MediaSignerLocal) SignConcatMP4(ctx context.Context, input io.ReadSeek startTime := time.Now() ctx, span := otel.Tracer("signer").Start(ctx, "SignMP4") defer span.End() - for _, ingredient := range ingredients { - _, err := iroh_streamplace.GetManifestAndCert(ingredient) - if err != nil { - return nil, err - } - } + // for _, ingredient := range ingredients { + // _, err := iroh_streamplace.GetManifestAndCert(c2patypes.NewReader(aqio.NewReadWriteSeeker(ingredient))) + // if err != nil { + // return nil, err + // } + // } // title := "livestream" mani := obj{ "title": "Livestream Clip", @@ -223,15 +224,12 @@ func (ms *MediaSignerLocal) SignConcatMP4(ctx context.Context, input io.ReadSeek } span.End() - bs, err := io.ReadAll(input) - if err != nil { - return nil, fmt.Errorf("failed to read input: %w", err) - } ctx, span = otel.Tracer("signer").Start(ctx, "SignMP4_Sign") rustCallbackSigner := &RustCallbackSigner{ Signer: ms.Signer, } - bs, err = iroh_streamplace.SignWithIngredients(string(manifestBs), bs, ms.Cert, ingredients, rustCallbackSigner) + rws := aqio.NewReadWriteSeeker([]byte{}) + err = iroh_streamplace.SignWithIngredients(string(manifestBs), c2patypes.NewReader(input), ms.Cert, ingredients, rustCallbackSigner, c2patypes.NewWriter(rws)) if err != nil { return nil, err } @@ -244,7 +242,7 @@ func (ms *MediaSignerLocal) SignConcatMP4(ctx context.Context, input io.ReadSeek } span.End() spmetrics.SigningDuration.WithLabelValues(ms.StreamerName).Observe(float64(time.Since(startTime).Milliseconds())) - return bs, nil + return rws.Bytes() } // don't call externally! this is used as a callback for the rust library diff --git a/pkg/media/segment_split.go b/pkg/media/segment_split.go index d4df26a3..35aa12aa 100644 --- a/pkg/media/segment_split.go +++ b/pkg/media/segment_split.go @@ -7,6 +7,7 @@ import ( "fmt" "sort" + "stream.place/streamplace/pkg/aqio" c2patypes "stream.place/streamplace/pkg/c2patypes" "stream.place/streamplace/pkg/iroh/generated/iroh_streamplace" "stream.place/streamplace/pkg/log" @@ -29,7 +30,7 @@ type ManifestAndMetadata struct { // split a signed concatenated mp4 into its constituent signed segments func SplitSegments(ctx context.Context, input []byte) ([]SplitSegment, error) { - manifestsStr, err := iroh_streamplace.GetManifests(input) + manifestsStr, err := iroh_streamplace.GetManifests(c2patypes.NewReader(aqio.NewReadWriteSeeker(input))) if err != nil { return nil, fmt.Errorf("failed to get manifests: %w", err) } @@ -68,7 +69,7 @@ func SplitSegments(ctx context.Context, input []byte) ([]SplitSegment, error) { if err != nil { return nil, fmt.Errorf("failed to segment file: %w", err) } - resignedSegs, err := iroh_streamplace.Resign(unsignedSegs, input, manifestStrs, certList) + resignedSegs, err := iroh_streamplace.Resign(unsignedSegs, c2patypes.NewReader(aqio.NewReadWriteSeeker(input)), manifestStrs, certList) if err != nil { return nil, fmt.Errorf("failed to resign segments: %w", err) } diff --git a/pkg/media/validate.go b/pkg/media/validate.go index 14308318..03303df7 100644 --- a/pkg/media/validate.go +++ b/pkg/media/validate.go @@ -12,6 +12,7 @@ import ( "github.com/bluesky-social/indigo/atproto/crypto" "go.opentelemetry.io/otel" + "stream.place/streamplace/pkg/aqio" "stream.place/streamplace/pkg/aqtime" c2patypes "stream.place/streamplace/pkg/c2patypes" "stream.place/streamplace/pkg/constants" @@ -209,7 +210,7 @@ type ValidationResult struct { // validate a signed mp4 file unto itself, ignoring whether this user is allowed and whatnot func ValidateMP4Media(ctx context.Context, buf []byte) (*ValidationResult, error) { var maniCert ManifestAndCert - maniStr, err := iroh_streamplace.GetManifestAndCert(buf) + maniStr, err := iroh_streamplace.GetManifestAndCert(c2patypes.NewReader(aqio.NewReadWriteSeeker(buf))) if err != nil { return nil, err } diff --git a/rust/iroh-streamplace/src/c2pa.rs b/rust/iroh-streamplace/src/c2pa.rs index 6b56c249..9437505d 100644 --- a/rust/iroh-streamplace/src/c2pa.rs +++ b/rust/iroh-streamplace/src/c2pa.rs @@ -10,20 +10,16 @@ use c2pa::jumbf_io; use c2pa::status_tracker::StatusTracker; use c2pa::store::Store; -use serde_json; +use crate::streams::Stream; +use crate::streams::StreamAdapter; -#[derive(Debug, thiserror::Error, uniffi::Error)] -#[uniffi(flat_error)] -pub enum SPError { - #[error("No certificate chain found")] - NoCertificateChainFound, - #[error("C2PA error: {0}")] - C2paError(String), -} +use serde_json; +use tracing::info; +use crate::error::SPError; #[uniffi::export] -pub fn get_manifest_and_cert(data: Vec) -> Result { - let reader = Reader::from_stream("video/mp4", Cursor::new(data)) +pub fn get_manifest_and_cert(data: &dyn Stream) -> Result { + let reader = Reader::from_stream("video/mp4", StreamAdapter::from(data)) .map_err(|e| SPError::C2paError(e.to_string()))?; if let Some(manifest) = reader.active_manifest() { @@ -113,7 +109,7 @@ enabled = true #[uniffi::export] pub fn sign( manifest: String, - data: Vec, + data: &dyn Stream, certs: Vec, gosigner: Arc, ) -> Result, SPError> { @@ -131,7 +127,7 @@ pub fn sign( let mut builder = Builder::from_json(&manifest).map_err(|e| SPError::C2paError(e.to_string()))?; let mut output = Vec::new(); - let mut input_cursor = Cursor::new(data); + let mut input_cursor = StreamAdapter::from(data); let mut output_cursor = Cursor::new(&mut output); builder .sign( @@ -145,8 +141,8 @@ pub fn sign( } #[uniffi::export] -pub fn get_manifests(data: Vec) -> Result { - let store = Reader::from_stream("video/mp4", Cursor::new(data)) +pub fn get_manifests(data: &dyn Stream) -> Result { + let store = Reader::from_stream("video/mp4", StreamAdapter::from(data)) .map_err(|e| SPError::C2paError(e.to_string()))?; let mut certs: std::collections::HashMap = std::collections::HashMap::new(); for (label, manifest) in store.manifests() { @@ -169,7 +165,7 @@ pub fn get_manifests(data: Vec) -> Result { #[uniffi::export] pub fn resign( unsigned_seg_data: Vec>, - signed_concat_data: Vec, + signed_concat_data: &dyn Stream, manifest_list: Vec, certs: Vec>, ) -> Result>, SPError> { @@ -177,7 +173,7 @@ pub fn resign( let combined_store = Store::from_stream( "video/mp4", - Cursor::new(signed_concat_data), + StreamAdapter::from(signed_concat_data), true, &mut validation_log, ) @@ -236,11 +232,12 @@ pub fn resign( #[uniffi::export] pub fn sign_with_ingredients( manifest: String, - data: Vec, + data: &dyn Stream, certs: Vec, ingredients: Vec>, gosigner: Arc, -) -> Result, SPError> { + output: &dyn Stream, +) -> Result<(), SPError> { Settings::from_toml(TOML_SETTINGS).map_err(|e| SPError::C2paError(e.to_string()))?; let callback_signer = CallbackSigner::new( move |_context: *const (), data: &[u8]| { @@ -253,15 +250,14 @@ pub fn sign_with_ingredients( ); let mut builder = Builder::from_json(&manifest).map_err(|e| SPError::C2paError(e.to_string()))?; - for ingredient in ingredients { + for (i, ingredient) in ingredients.into_iter().enumerate() { let mut cursor = Cursor::new(ingredient); let ingredient = Ingredient::from_stream("video/mp4", &mut cursor) .map_err(|e| SPError::C2paError(e.to_string()))?; builder.add_ingredient(ingredient); } - let mut output = Vec::new(); - let mut input_cursor = Cursor::new(data); - let mut output_cursor = Cursor::new(&mut output); + let mut input_cursor = StreamAdapter::from(data); + let mut output_cursor = StreamAdapter::from(output); builder .sign( &callback_signer, @@ -270,5 +266,5 @@ pub fn sign_with_ingredients( &mut output_cursor, ) .map_err(|e| SPError::C2paError(e.to_string()))?; - Ok(output) + Ok(()) } diff --git a/rust/iroh-streamplace/src/error.rs b/rust/iroh-streamplace/src/error.rs new file mode 100644 index 00000000..0e5d2e49 --- /dev/null +++ b/rust/iroh-streamplace/src/error.rs @@ -0,0 +1,10 @@ +#[derive(Debug, thiserror::Error, uniffi::Error)] +#[uniffi(flat_error)] +pub enum SPError { + #[error("No certificate chain found")] + NoCertificateChainFound, + #[error("C2PA error: {0}")] + C2paError(String), + #[error("IO Error: {0}")] + IOError(String), +} diff --git a/rust/iroh-streamplace/src/lib.rs b/rust/iroh-streamplace/src/lib.rs index 9f335aef..709cc7a4 100644 --- a/rust/iroh-streamplace/src/lib.rs +++ b/rust/iroh-streamplace/src/lib.rs @@ -1,8 +1,10 @@ uniffi::setup_scaffolding!(); pub mod c2pa; +pub mod error; pub mod node_addr; pub mod public_key; +pub mod streams; use std::sync::{LazyLock, Once}; @@ -11,6 +13,9 @@ pub use db::*; #[cfg(test)] mod tests; +#[cfg(test)] +mod test_stream; + /// Lazily initialized Tokio runtime for use in uniffi methods that need a runtime. static RUNTIME: LazyLock = LazyLock::new(|| tokio::runtime::Runtime::new().unwrap()); diff --git a/rust/iroh-streamplace/src/streams.rs b/rust/iroh-streamplace/src/streams.rs index e69de29b..e8e1d413 100644 --- a/rust/iroh-streamplace/src/streams.rs +++ b/rust/iroh-streamplace/src/streams.rs @@ -0,0 +1,167 @@ +// Copyright 2023 Adobe. All rights reserved. +// This file is licensed to you under the Apache License, +// Version 2.0 (http://www.apache.org/licenses/LICENSE-2.0) +// or the MIT license (http://opensource.org/licenses/MIT), +// at your option. + +// Unless required by applicable law or agreed to in writing, +// this software is distributed on an "AS IS" BASIS, WITHOUT +// WARRANTIES OR REPRESENTATIONS OF ANY KIND, either express or +// implied. See the LICENSE-MIT and LICENSE-APACHE files for the +// specific language governing permissions and limitations under +// each license. + +use crate::error::SPError; +use std::io::{Read, Seek, SeekFrom, Write}; +use std::sync::Arc; + +// #[repr(C)] +// #[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] +// pub enum SeekMode { +// Start = 0, +// End = 1, +// Current = 2, +// } + +/// This allows for a callback stream over the Uniffi interface. +/// Implement these stream functions in the foreign language +/// and this will provide Rust Stream trait implementations +/// This is necessary since the Rust traits cannot be implemented directly +/// as uniffi callbacks +#[uniffi::export(with_foreign)] +pub trait Stream: Send + Sync { + /// Read a stream of bytes from the stream + fn read_stream(&self, length: u64) -> Result, SPError>; + /// Seek to a position in the stream + fn seek_stream(&self, pos: i64, mode: u64) -> Result; + /// Write a stream of bytes to the stream + fn write_stream(&self, data: Vec) -> Result; +} + +impl Stream for Arc { + fn read_stream(&self, length: u64) -> Result, SPError> { + (**self).read_stream(length) + } + + fn seek_stream(&self, pos: i64, mode: u64) -> Result { + (**self).seek_stream(pos, mode) + } + + fn write_stream(&self, data: Vec) -> Result { + (**self).write_stream(data) + } +} + +impl AsMut for dyn Stream { + fn as_mut(&mut self) -> &mut Self { + self + } +} + +pub struct StreamAdapter<'a> { + pub stream: &'a mut dyn Stream, +} + +impl<'a> StreamAdapter<'a> { + pub fn from_stream_mut(stream: &'a mut dyn Stream) -> Self { + Self { stream } + } +} + +impl<'a> From<&'a dyn Stream> for StreamAdapter<'a> { + #[allow(invalid_reference_casting)] + fn from(stream: &'a dyn Stream) -> Self { + let stream = &*stream as *const dyn Stream as *mut dyn Stream; + let stream = unsafe { &mut *stream }; + Self { stream } + } +} + +// impl<'a> c2pa::CAIRead for StreamAdapter<'a> {} + +// impl<'a> c2pa::CAIReadWrite for StreamAdapter<'a> {} + +impl<'a> Read for StreamAdapter<'a> { + fn read(&mut self, buf: &mut [u8]) -> std::io::Result { + let mut bytes = self + .stream + .read_stream(buf.len() as u64) + .map_err(|e| std::io::Error::new(std::io::ErrorKind::Other, e))?; + let len = bytes.len(); + buf.iter_mut().zip(bytes.drain(..)).for_each(|(dest, src)| { + *dest = src; + }); + //println!("read: {:?}", len); + Ok(len) + } +} + +impl<'a> Seek for StreamAdapter<'a> { + fn seek(&mut self, pos: std::io::SeekFrom) -> std::io::Result { + let (pos, mode) = match pos { + SeekFrom::Current(pos) => (pos, 2), + SeekFrom::Start(pos) => (pos as i64, 0), + SeekFrom::End(pos) => (pos, 1), + }; + //println!("Stream Seek {}", pos); + self.stream + .seek_stream(pos, mode) + .map_err(|e| std::io::Error::new(std::io::ErrorKind::Other, e)) + } +} + +impl<'a> Write for StreamAdapter<'a> { + fn write(&mut self, buf: &[u8]) -> std::io::Result { + let len = self + .stream + .write_stream(buf.to_vec()) + .map_err(|e| std::io::Error::new(std::io::ErrorKind::Other, e))?; + Ok(len as usize) + } + + fn flush(&mut self) -> std::io::Result<()> { + Ok(()) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + use crate::test_stream::TestStream; + + #[test] + fn test_stream_read() { + let mut test = TestStream::from_memory(vec![0, 1, 2, 3, 4, 5, 6, 7, 8, 9]); + let mut stream = StreamAdapter::from_stream_mut(&mut test); + let mut buf = [0u8; 5]; + let len = stream.read(&mut buf).unwrap(); + assert_eq!(len, 5); + assert_eq!(buf, [0, 1, 2, 3, 4]); + } + + #[test] + fn test_stream_seek() { + let mut test = TestStream::from_memory(vec![0, 1, 2, 3, 4, 5, 6, 7, 8, 9]); + let mut stream = StreamAdapter { stream: &mut test }; + let pos = stream.seek(SeekFrom::Start(5)).unwrap(); + assert_eq!(pos, 5); + let mut buf = [0u8; 5]; + let len = stream.read(&mut buf).unwrap(); + assert_eq!(len, 5); + assert_eq!(buf, [5, 6, 7, 8, 9]); + } + + #[test] + fn test_stream_write() { + let mut test = TestStream::new(); + let mut stream = StreamAdapter { stream: &mut test }; + let len = stream.write(&[0, 1, 2, 3, 4]).unwrap(); + assert_eq!(len, 5); + stream.seek(SeekFrom::Start(0)).unwrap(); + let mut buf = [0u8; 5]; + let len = stream.read(&mut buf).unwrap(); + assert_eq!(len, 5); + assert_eq!(buf, [0, 1, 2, 3, 4]); + } +} diff --git a/rust/iroh-streamplace/src/test_stream.rs b/rust/iroh-streamplace/src/test_stream.rs new file mode 100644 index 00000000..2ca30bcb --- /dev/null +++ b/rust/iroh-streamplace/src/test_stream.rs @@ -0,0 +1,80 @@ +// Copyright 2023 Adobe. All rights reserved. +// This file is licensed to you under the Apache License, +// Version 2.0 (http://www.apache.org/licenses/LICENSE-2.0) +// or the MIT license (http://opensource.org/licenses/MIT), +// at your option. + +// Unless required by applicable law or agreed to in writing, +// this software is distributed on an "AS IS" BASIS, WITHOUT +// WARRANTIES OR REPRESENTATIONS OF ANY KIND, either express or +// implied. See the LICENSE-MIT and LICENSE-APACHE files for the +// specific language governing permissions and limitations under +// each license. + +use std::io::{Read, Seek, SeekFrom, Write}; +use std::sync::RwLock; + +use crate::error::SPError; +use crate::streams::Stream; +use std::io::Cursor; + +pub struct TestStream { + stream: RwLock>>, +} + +impl TestStream { + pub fn new() -> Self { + Self { + stream: RwLock::new(Cursor::new(Vec::new())), + } + } + pub fn from_memory(data: Vec) -> Self { + Self { + stream: RwLock::new(Cursor::new(data)), + } + } +} + +impl Stream for TestStream { + fn read_stream(&self, length: u64) -> Result, SPError> { + if let Ok(mut stream) = RwLock::write(&self.stream) { + let mut data = vec![0u8; length as usize]; + let bytes_read = stream + .read(&mut data) + .map_err(|e| SPError::IOError(e.to_string()))?; + data.truncate(bytes_read); + //println!("read_stream: {:?}, pos {:?}", data.len(), (*stream).position()); + Ok(data) + } else { + Err(SPError::IOError("RwLock".to_string()))? + } + } + + fn seek_stream(&self, pos: i64, mode: u64) -> Result { + if let Ok(mut stream) = RwLock::write(&self.stream) { + //stream.seek(SeekFrom::Start(pos as u64)).map_err(|e| StreamError::Io{ reason: e.to_string()})?; + let whence = match mode { + 0 => SeekFrom::Start(pos as u64), + 1 => SeekFrom::End(pos as i64), + 2 => SeekFrom::Current(pos as i64), + 3_u64..=u64::MAX => unimplemented!(), + }; + Ok(stream + .seek(whence) + .map_err(|e| SPError::IOError(e.to_string()))?) + } else { + Err(SPError::IOError("RwLock".to_string())) + } + } + + fn write_stream(&self, data: Vec) -> Result { + if let Ok(mut stream) = RwLock::write(&self.stream) { + let len = stream + .write(&data) + .map_err(|e| SPError::IOError(e.to_string()))?; + Ok(len as u64) + } else { + Err(SPError::IOError("RwLock".to_string()))? + } + } +}