Skip to content

Commit 1d50eea

Browse files
committed
add streaming
1 parent 7455f98 commit 1d50eea

6 files changed

Lines changed: 517 additions & 0 deletions

File tree

‎lib/iq_rdf.rb‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,7 @@
2626
require 'iq_rdf/collection'
2727
require 'iq_rdf/predicate_namespace'
2828
require 'iq_rdf/document'
29+
require 'iq_rdf/streaming_writer'
2930

3031
require 'builder'
3132

‎lib/iq_rdf/document.rb‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,10 @@
1515
module IqRdf
1616
class Document
1717

18+
def self.stream(io, format, **opts, &block)
19+
IqRdf::StreamingWriter.open(io, format, **opts, &block)
20+
end
21+
1822
def initialize(default_namespace_uri_prefix = nil, *args)
1923
options = args.last.is_a?(::Hash) ? args.pop : {}
2024
raise ArgumentError, "If given, parameter :lang has to be a Symbol" unless options[:lang].nil? || options[:lang].is_a?(Symbol)

‎lib/iq_rdf/streaming_writer.rb‎

Lines changed: 202 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,202 @@
1+
module IqRdf
2+
class StreamingWriter
3+
4+
FORMATS = %w[ttl nt xml].freeze
5+
6+
def initialize(io, format, default_namespace: nil, lang: nil, config: {})
7+
@io = io
8+
@format = format.to_s
9+
raise ArgumentError, "Unknown format '#{format}'. Must be one of: #{FORMATS.join(', ')}" unless FORMATS.include?(@format)
10+
@document_language = lang
11+
@config = config
12+
@namespaces = {}
13+
@header_written = false
14+
@blank_nodes = {}
15+
16+
register_namespace(:rdf, URI.parse("http://www.w3.org/1999/02/22-rdf-syntax-ns#"))
17+
namespaces(default: default_namespace) if default_namespace
18+
end
19+
20+
def namespaces(namespaces)
21+
raise ArgumentError, "Parameter 'namespaces' has to be a hash" unless namespaces.is_a?(Hash)
22+
namespaces.each do |name, uri_prefix|
23+
uri_prefix = ::URI.parse(uri_prefix)
24+
raise ArgumentError, "Parameter 'namespaces' must be in the form {Symbol => URIString, ...}" unless name.is_a?(Symbol)
25+
register_namespace(name, uri_prefix)
26+
end
27+
self
28+
end
29+
30+
def <<(node)
31+
return if node.nil?
32+
raise ArgumentError, "Node must be an IqRdf::Uri and a Subject!" unless node.is_a?(IqRdf::Uri) && node.is_subject?
33+
write_header unless @header_written
34+
serialize_node(node)
35+
end
36+
37+
def close
38+
write_xml_footer if @format == 'xml' && @header_written
39+
end
40+
41+
def self.open(io, format, **opts)
42+
writer = new(io, format, **opts)
43+
yield writer
44+
ensure
45+
writer&.close
46+
end
47+
48+
private
49+
50+
def register_namespace(name, uri_prefix)
51+
@namespaces[name] = IqRdf::Namespace.create(name, uri_prefix)
52+
end
53+
54+
def write_header
55+
case @format
56+
when 'ttl' then write_turtle_header
57+
when 'xml' then write_xml_header
58+
end
59+
@header_written = true
60+
end
61+
62+
def serialize_node(subject)
63+
case @format
64+
when 'ttl' then serialize_turtle(subject)
65+
when 'nt' then serialize_ntriples(subject)
66+
when 'xml' then serialize_xml(subject)
67+
end
68+
end
69+
70+
# --- Turtle ---
71+
72+
def write_turtle_header
73+
@namespaces.values.sort_by(&:turtle_token).each do |namespace|
74+
@io.write("@prefix #{namespace.turtle_token}: <#{namespace.uri_prefix}>.\n")
75+
end
76+
@io.write("\n")
77+
end
78+
79+
def serialize_turtle(subject)
80+
pref = subject.to_s
81+
indent = "".ljust(pref.length)
82+
83+
if subject.rdf_type
84+
@io.write("#{pref} a #{subject.rdf_type}")
85+
pref = ";\n" + indent
86+
end
87+
88+
subject.nodes.each do |predicate|
89+
objects = predicate.nodes.map { |object|
90+
object.to_s(indent: indent, lang: predicate.lang || subject.lang || @document_language)
91+
}.join(", ")
92+
@io.write("#{pref} #{predicate} #{objects}")
93+
pref = ";\n" + indent
94+
end
95+
@io.write(".\n")
96+
@io.write("\n") if @config[:empty_line_between_triples]
97+
end
98+
99+
# --- N-Triples ---
100+
101+
def serialize_ntriples(sbj)
102+
nt_process_subject(sbj) do |triple|
103+
nt_process_blank_nodes(triple, sbj)
104+
end
105+
end
106+
107+
def nt_process_subject(sbj, &block)
108+
rdf_type = IqRdf::Rdf::build_uri("type")
109+
110+
if (sbj.rdf_type rescue false)
111+
lang = sbj.lang || @document_language
112+
nt_write_triple([sbj, rdf_type, sbj.rdf_type], lang)
113+
end
114+
115+
sbj.nodes.each do |prd|
116+
lang = prd.lang || sbj.lang || @document_language
117+
prd.nodes.each do |obj|
118+
triple = [sbj, prd, obj]
119+
nt_write_triple(triple, lang)
120+
block.call(triple) if block
121+
end
122+
end
123+
end
124+
125+
def nt_process_blank_nodes(triple, current_res)
126+
sbj, _prd, obj = triple
127+
[sbj, obj].select { |res| res.is_a?(IqRdf::BlankNode) && res != current_res }.each do |res|
128+
nt_process_subject(res) do |inner_triple|
129+
nt_process_blank_nodes(inner_triple, res)
130+
end
131+
end
132+
end
133+
134+
def nt_write_triple(triple, lang)
135+
parts = triple.map { |res| nt_resource(res, lang) }
136+
@io.write("#{parts.join(' ')} .\n")
137+
end
138+
139+
def nt_resource(res, lang)
140+
if res.is_a?(IqRdf::Literal)
141+
res.to_ntriples(lang)
142+
elsif res.is_a?(IqRdf::BlankNode)
143+
nt_blank_node(res)
144+
elsif res.is_a?(IqRdf::Collection)
145+
nt_collection(res)
146+
else
147+
"<#{res.full_uri}>"
148+
end
149+
end
150+
151+
def nt_blank_node(res)
152+
@blank_nodes[res] ||= @blank_nodes.size + 1
153+
"_:b#{@blank_nodes[res]}"
154+
end
155+
156+
def nt_collection(res)
157+
nt_blank_node(res) # register collection object (matches original side-effect)
158+
list = IqRdf::BlankNode.new
159+
sublist = list
160+
total = res.elements.length
161+
res.elements.each_with_index do |current_element, i|
162+
last = i + 1 == total
163+
sublist::rdf.build_predicate("type", IqRdf::Rdf::build_uri("List"))
164+
sublist::rdf.first(current_element)
165+
if last
166+
sublist::rdf.rest(IqRdf::Rdf::build_uri("nil"))
167+
else
168+
new_sublist = IqRdf::BlankNode.new
169+
sublist::rdf.rest(new_sublist)
170+
end
171+
nt_process_subject(sublist) { |triple| nt_process_blank_nodes(triple, sublist) }
172+
sublist = new_sublist unless last
173+
end
174+
nt_blank_node(list)
175+
end
176+
177+
# --- XML ---
178+
179+
def write_xml_header
180+
@xml = Builder::XmlMarkup.new(target: @io, indent: 2)
181+
@xml.instruct!
182+
opts = {}
183+
@namespaces.values.each do |namespace|
184+
opts[namespace.token == :default ? "xmlns" : "xmlns:#{namespace.token}"] = namespace.uri_prefix
185+
end
186+
opts["xml:lang"] = @document_language if @document_language
187+
# Fiber pausiert den Builder-Block nach dem öffnenden Tag, so dass
188+
# Builder die Einrückungstiefe korrekt trackt während Nodes gestreamt werden.
189+
@xml_fiber = Fiber.new { @xml.rdf(:RDF, opts) { Fiber.yield } }
190+
@xml_fiber.resume
191+
end
192+
193+
def serialize_xml(node)
194+
node.build_xml(@xml)
195+
end
196+
197+
def write_xml_footer
198+
@xml_fiber.resume if @xml_fiber&.alive?
199+
end
200+
201+
end
202+
end

‎test/streaming_ntriples_test.rb‎

Lines changed: 115 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,115 @@
1+
# -*- encoding : utf-8 -*-
2+
$LOAD_PATH << File.dirname(__FILE__)
3+
4+
require 'test_helper'
5+
require 'stringio'
6+
7+
class StreamingNTriplesTest < Minitest::Test
8+
9+
def stream(default_namespace: 'http://www.test.de/', lang: nil, &block)
10+
io = StringIO.new
11+
IqRdf::Document.stream(io, :nt, default_namespace: default_namespace, lang: lang, &block)
12+
io.string
13+
end
14+
15+
def test_basics
16+
result = stream(lang: :de) do |doc|
17+
doc.namespaces foaf: 'http://xmlns.com/foaf/0.1/'
18+
doc << IqRdf::testemann do |t|
19+
t.Foaf::knows(IqRdf::testefrau)
20+
t.Foaf.nick("Testy")
21+
t.Foaf.lastname("Testemann", :lang => :none)
22+
end
23+
end
24+
25+
assert_equal(<<~RDF.strip, result.strip)
26+
<http://www.test.de/testemann> <http://xmlns.com/foaf/0.1/knows> <http://www.test.de/testefrau> .
27+
<http://www.test.de/testemann> <http://xmlns.com/foaf/0.1/nick> "Testy"@de .
28+
<http://www.test.de/testemann> <http://xmlns.com/foaf/0.1/lastname> "Testemann" .
29+
RDF
30+
assert result.end_with?("\n"), "should end with trailing newline"
31+
end
32+
33+
def test_matches_document_output
34+
nodes = lambda do |doc|
35+
doc.namespaces skos: 'http://www.w3.org/2008/05/skos#',
36+
foaf: 'http://xmlns.com/foaf/0.1/',
37+
upb: 'http://www.upb.de/'
38+
doc << IqRdf::testemann.myCustomNote("This is an example", :lang => :en)
39+
doc << IqRdf::testemann(IqRdf::Foaf::build_uri("Person")).Foaf::name("Heinz Peter Testemann", :lang => :none)
40+
doc << IqRdf::testemann.Foaf::knows(IqRdf::testefrau)
41+
doc << IqRdf::testemann.Foaf::nick("Crash test dummy")
42+
end
43+
44+
document = IqRdf::Document.new('http://www.umweltprobenbank.de/', :lang => :de)
45+
nodes.call(document)
46+
expected = document.to_ntriples
47+
48+
result = stream(default_namespace: 'http://www.umweltprobenbank.de/', lang: :de) do |doc|
49+
nodes.call(doc)
50+
end
51+
52+
assert_equal expected, result
53+
end
54+
55+
def test_blank_nodes
56+
result = stream do |doc|
57+
doc << IqRdf::testnode.test32 do |blank_node|
58+
blank_node.title("dies ist ein test")
59+
blank_node.build_predicate(:test, "Another test")
60+
blank_node.sub do |subnode|
61+
subnode.title("blubb")
62+
end
63+
end
64+
end
65+
66+
assert_equal(<<~RDF.strip, result.strip)
67+
<http://www.test.de/testnode> <http://www.test.de/test32> _:b1 .
68+
_:b1 <http://www.test.de/title> "dies ist ein test" .
69+
_:b1 <http://www.test.de/test> "Another test" .
70+
_:b1 <http://www.test.de/sub> _:b2 .
71+
_:b2 <http://www.test.de/title> "blubb" .
72+
RDF
73+
end
74+
75+
def test_blank_nodes_numbered_across_nodes
76+
result = stream do |doc|
77+
doc << IqRdf::node1.pred1 do |blank_node|
78+
blank_node.title("first")
79+
end
80+
doc << IqRdf::node2.pred2 do |blank_node|
81+
blank_node.title("second")
82+
end
83+
end
84+
85+
# blank node counter must not reset between << calls:
86+
# second blank node gets _:b2, not _:b1 again
87+
assert_match(/_:b2 <.*> "second"/, result)
88+
refute_match(/_:b1 <.*> "second"/, result)
89+
end
90+
91+
def test_collections
92+
result = stream(default_namespace: 'http://test.de/') do |doc|
93+
doc << IqRdf::testemann.testIt([IqRdf::hello, IqRdf::goodbye, "bla"])
94+
end
95+
96+
assert_match(/<http:\/\/test.de\/testemann> <http:\/\/test.de\/testIt> _:b/, result)
97+
assert_match(/rdf-syntax-ns#List/, result)
98+
assert_match(/rdf-syntax-ns#first/, result)
99+
assert_match(/rdf-syntax-ns#rest/, result)
100+
end
101+
102+
def test_nil_node_is_ignored
103+
result = stream { |doc| doc << nil }
104+
assert_equal "", result
105+
end
106+
107+
def test_no_header_written
108+
io = StringIO.new
109+
IqRdf::Document.stream(io, :nt, default_namespace: 'http://www.test.de/') do |doc|
110+
doc << IqRdf::testemann.title("hello")
111+
end
112+
refute_match(/@prefix/, io.string, "NT format should have no @prefix header")
113+
end
114+
115+
end

0 commit comments

Comments
 (0)